macOS - TECH PLAY - TECH PLAY

TECH PLAY

macOS

イベント

マガジン

該当するコンテンツが見つかりませんでした

技術ブログ

チューリングの MLOps エンジニアの岩政です。 以前、Databricks の Declarative Automation Bundles を用いた機械学習データセット作成基盤の構築をまとめました。この記事はその続きにあたるので、まだご覧になっていない方は先に読んでいただけると嬉しいです。 https://zenn.dev/turing_motors/articles/46b5560dce0e3a この中で「4.5.2 Python のライブラリの依存解決を高速にする」を紹介しました。 Databricks の Serverless compute では、依存パッケージのインストー
こんにちは。データサイエンティストの白井です。 今日は、「PPTXやPDFをLLMに読ませるとき、入力をどう刻むと内容理解の精度が変わるのか」を実測した話を紹介します。 1枚ずつ渡すのか、数枚まとめるのか、PDFで丸ごと渡すのか。この違いだけで、精度・処理時間・コストがどう動くかを揃えた条件で比較しました。 先に本記事のポイントを書いておきます。 資料全体を統合した最終的なまとめは、どの刻み方でもほとんど変わりませんでした。 差が出たのは 1ページごとの内容理解の細部 でした。注記や例外条件を拾えるか、グラフの軸や目盛りを正しく読めるか、併記された数値の帰属先を取り違えないか。この記事は、その細部の話です。 はじめに 検証の設計 対象資料と出典 比較する3方式 統制した条件 読み取るべき内容の整理 結果 3方式の総合比較 増えた分に何が入っていたか スライド9 スライド7 使い分けの指針 おまけ:PPTXをPNG化すると画像がズレる おわりに 出典・注意事項 はじめに 私は日々の業務で様々なパワーポイントの資料(PPTX)を目にします。 LLMを業務に活用する場合、LLMにPPTXの内容を理解してもらうことが必須です。 そして、そのコンテキストの上で、LLMと一緒に作業を進めることになります。 社内の提案書や行政資料には、1枚の中にグラフ・注記・用語の定義・分析コメントがまとめて載っているスライドがよくあります。 今回の検証で使った資料の1枚がこれです。 ※出典: RESAS Portal - 地域課題分析ナビゲーション テーマ① 地域の人口減少対策 - (画像化して掲載) 左に散布図、右に分析の視点、下に用語の説明。 そして散布図の上には「鹿児島県薩摩川内市」「全国平均」「説明変数(X軸)38.54」「目的変数(Y軸)78.54」が縦に並んでいます。 この38.54と78.54が どちらに係る値なのか (当該市の値なのか、全国平均の値なのか)は、行間を読まないと分かりません。 実はこの箇所、後述するゴールドデータを人手で作ったときに私自身も一度間違えました。 担当者が分析手順をそのままたどれるように、1枚で「使うデータ」「分析の視点」「用語の定義」まで完結させる構成なのだろうと推測しています。人が手順書として読むには親切ですが、LLMに渡す側から見ると、1枚あたりに読み取るべき情報が多い資料ということになります。 こういう資料をLLMに渡すとき、地味に悩むのが PPTXの渡し方 です。 PPTXが標準オブジェクトだけで作られていればLLMは正確に読めますが、画像貼り付けが使われている場合はvision機能が必要になります。私の周りの大抵のPPTXには画像貼り付けが使われているため、私はいつも全体をPNGにして渡していました。 しかし、ふと、この渡し方で十分だろうか?と疑問が湧き、以下3方式を検証してみようと思いました。 全部まとめてPDFで1回投げるのか 1ページずつ投げるのか 数ページずつブロックで投げるのか 「1枚ずつが丁寧そうだが、リクエスト数が増えて時間もコストも膨らむ」というのは感覚的に分かります。 では、実際にどれだけ細かく読み取れるようになって、どれだけ高くつくのか。 同じ資料・同じプロンプト・同じモデルで3方式を各3試行ずつ回して、その数字を出しました。 方式 リクエスト数 1枚あたり facts件数 処理時間 コスト ① PDF一括 2 2.14 195秒 $0.29 ② 3枚ずつ 9 4.33 573秒 $0.74 ③ 1枚ずつ 21 5.98 816秒 $1.00 1枚ずつは、1スライドから読み取る事実がPDF一括の約2.8倍。ただし時間は約4倍、コストは約3.5倍です。 そして本当の論点は、この「約2.8倍」が中身を伴った差なのかどうかです。件数が多いだけなら意味がありません。記事の後半では、増えた分に何が入っていたかを実際の出力で見ていきます。 検証の設計 対象資料と出典 使ったのは、RESAS Portalで公開されている「地域課題分析ナビゲーション テーマ① 地域の人口減少対策」というPPTXです。 自治体が人口減少の要因を分析する手順を、RESASの実際のグラフを使って解説する資料で、鹿児島県薩摩川内市を例にしています。 全27枚のうち、分析パートである スライド5〜23の19枚 を入力対象にしました。 なお、本記事内でのスライドの利用にあたっては、RESAS Portalのサイトポリシー(政府標準利用規約準拠)を確認したうえで、記事化の可否を問い合わせ窓口にメールで確認し、許諾をいただいています。 出典表記と加工の明示、免責事項は記事末尾にまとめました。 比較する3方式 比較したのは、入力の刻み方だけが違う3方式です。 モデルは全方式で Claude Sonnet 5 を使いました。 ① PDF一括 ② 3枚ずつ ③ 1枚ずつ 入力単位 19枚をラスタライズしたPDF 1本 PNG 3枚ずつ(7ブロック) PNG 1枚ずつ(19枚) 解析リクエスト 1回 7回 19回 最終統合 1回 1回 1回 どの方式も処理は2段階です。 スライド解析 :渡された画像の主メッセージ・重要な数値・注記・読み取れない箇所をJSONで出させる 最終統合 :1の結果を全部集めて、資料全体の人口減少要因・優先課題・追加調査事項を整理させる 冒頭に書いたとおり、2の最終統合の出来は3方式でほとんど変わりませんでした。以降で見ていくのは1のスライド解析についてのみです(統合ステップも時間とコストは消費するので、計測には含めています)。 ポイントは、 利用するプロンプトを3方式で完全に同一にしている ことです。 「PDFで渡すときはスライド番号を意識してね」といった補助を足すと、比較しているものが「入力単位の違い」から「プロンプトの違い」に変わってしまうので、あえて共通のまま使いました。 スライド番号のヒントも与えていません。元資料の右下にページ番号が印字されているので、PDF一括でもPNG単体でも、画像から得られる情報量は同じという状態にしてあります。 統制した条件 「入力の刻み方」以外の差は、できるだけ合わせるようにしています。 同じPNG画像を使う :3方式すべて同一のレンダリング結果。PDFも、そのPNGを結合してラスタライズしたもの 同じプロンプト・同じ出力スキーマ :解析・統合ともに共通。JSON Schema違反時は1回だけ修復リトライ、という条件も統一 Prompt Cachingは無効 :キャッシュが効くと「刻み方によるコスト差」が見えなくなるため Batch APIは使わない・並列実行しない :処理時間を素直に比較するため、逐次実行 temperature等は指定しない :モデルのデフォルト 各方式3試行 :LLMの出力はブレるので、1回では判断しない 実行順はランダム化 :時間帯によるAPIの混み具合が特定方式に偏らないように 最大出力トークンも揃えていますが、ここは1点だけ例外があります。 PDF一括方式は1リクエストで19枚分の解析結果を出す必要があるため、通常の6,000では確実に途中で切れます。 そのためPDF一括の解析のみ32,000を別枠で設定しました。 この統一した条件で、以下の内容をJSONフォーマットとして出力させました。 main_message: スライドのメインメッセージ facts: 読み取った事実 important_numbers: 重要指標 conditions_and_notes: 条件や備考 unclear_points: 不明点 利用したプロンプト(クリックで展開) system_prompt # 役割と目的 あなたは、日本語の行政資料(RESAS Portal「地域課題分析ナビゲーション テーマ① 地域の人口減少対策」)の内容理解を支援するアシスタントです。 この資料は、特定自治体の政策を決定するための完成済み提言書ではなく、RESAS等を用いた地域課題分析の進め方・分析例を示す教材です。あなたの役割は、資料の内容を正確に把握し、人口減少への対応として優先的に検討すべき課題とネクストアクションの叩き台を整理することです。**政策を決定すること自体は行いません。** # 厳守事項 1. **資料に書かれていない事実を補完しない。** 資料から読み取れないことは、推測せず「不確実」「読み取れない」として明示する。 2. **数値・期間・対象地域・年齢区分を正確に区別する。** 自然増減(出生・死亡)と社会増減(転入・転出)を混同しない。 3. **すべての主張に根拠スライド番号を付与する。** 根拠のない主張は出力しない。 4. **出力は「政策決定」ではない。**(後略:政策を断定的に推奨しない) 5. **人口移動や出生に関する分析を、特定の年代・性別・個人の責任として表現しない。** 構造的要因として扱う。 6. **出力は指定されたJSON Schemaに厳密に従う。** # 不確実性の扱い - 数値や文言が読み取れない場合は、無理に推測せず`unclear_points`等の該当フィールドへ記録する。 - 資料内の記載が矛盾する場合は、矛盾として明示し、どちらか一方を勝手に採用しない。 - 資料だけでは判断できない事項は、`limitations`や`cannot_be_decided_from_material`に明示する。 user prompt # スライド/ブロック解析 入力として渡されたスライド(画像、または画像とテキストの組み合わせ)を解析してください。 ## 手順 1. 入力された各スライドについて、主メッセージ・重要な数値・条件や注記・読み取れない箇所を整理する。 2. 複数スライドが同時に入力された場合は、スライド間で関連する内容があれば`cross_slide_findings`に記録する。 3. 資料全体の主張につながりそうな課題の候補があれば`candidate_issues`に記録する(この時点では確定させない)。 4. 疑問点・追加確認したい点は`open_questions`に記録する。 5. **同一スライドについてスライド全体画像に加えて部分拡大画像(crop)が入力される場合**、入力順は「スライド全体」→「左上」→「右上」→「左下」→「右下」(相互に10〜15%重なる)である。cropは細部(小さな数値・注記等)の確認に使い、全体像の把握はスライド全体画像を優先する。 6. **画像に加えて、同一スライドのPowerPointから抽出したテキストが入力される場合**、画像から読み取れる内容とテキストの内容を突き合わせ、食い違いがあれば推測で解決せず`conditions_and_notes`にその旨(どちらの内容がどう異なるか)を記録する。 ## 出力形式 以下のJSON Schemaに厳密に従うJSONのみを出力してください。前置き・後書き・Markdownのコードフェンスは不要です。 { "processed_slides": [5, 6, 7], "slide_findings": [ { "slide_number": 5, "main_message": "", "facts": [ { "statement": "", "evidence": "", "confidence": "high|medium|low" } ], "important_numbers": [], "conditions_and_notes": [], "unclear_points": [] } ], "cross_slide_findings": [], "candidate_issues": [], "open_questions": [] } important_numbers: 資料中の重要な数値を、単位・対象期間・対象地域が分かる形の文字列で記録する(例: "2065年に老年人口割合が約4割へ上昇")。 conditions_and_notes: グラフの前提条件、注記、例外条件。 unclear_points: 読み取れない・判断できない箇所。推測しない。 cross_slide_findings・candidate_issues・open_questionsの各要素は、オブジェクトではなく1つの文字列にしてください。根拠となるスライド番号を示したい場合は、文字列の中に含めてください。 読み取るべき内容の整理 出力の良し悪しを自分で判断できるように、スライド5〜23の 19枚すべてについて、読み取るべき内容を人手で整理しました (以下、ゴールドデータと呼びます)。 各スライドについて、以下をJSONで記述しています。 main_message :そのスライドの主メッセージ key_numbers :重要な数値・定義・対象期間 graph_trends :グラフの増減・比較関係 notes_and_exceptions :注記・前提条件・例外 easily_misread_points :読み違えやすい点 例えば、以下のスライド22は、 ゴールドデータとして以下のjsonを作成しています。 { "slide_number": 22, "main_message": "婚姻件数は令和2年に一時的に減少しているが大きな変動はなく、夫婦の初婚平均年齢も夫30歳前後・妻28歳前後で安定して推移している。", "key_numbers": [ "婚姻件数: 平成25年416件〜令和2年354件(令和元年442件がピーク)", "夫の初婚平均年齢29.2〜30.6歳(最小:平成27年29.2歳、最大:令和元年30.6歳)", "妻の初婚平均年齢28.2〜29.3歳(最小:令和2年28.2歳、最大:令和元年29.3歳)" ], "graph_trends": [ "婚姻件数は416→373→396→392→434→442→354件と増減を繰り返しており単調増加ではない(平成26年・平成28年に一時的な減少)。令和元年442件がピークで、令和2年に354件へ減少", "夫・妻の平均初婚年齢はともに横ばい(変動幅は夫1.4歳・妻1.1歳の範囲内)" ], "notes_and_exceptions": ["出典は鹿児島県「人口動態統計調査」"], "easily_misread_points": [ "左軸(婚姻件数)と右軸(平均年齢)のスケールが異なる二軸グラフのため、棒グラフと折れ線グラフの変動幅を単純比較すると誤解を招く", "夫-平均と妻-平均の2本の折れ線のデータラベルが近接しており、夫の値(令和2年29.4歳)を妻の値と取り違えやすい(本ゴールドデータの初版が実際にこの誤りを犯した)", "分析の視点欄に「婚姻件数は令和2年に一時的に減少しているものの、大きな変動は見られていない」とあるため、実データが増減を繰り返している事実を「単調増加後の減少」と単純化しやすい" ] }, このゴールドはまずLLMに作成させ、その後筆者が結果をレビューし、記載内容が正しいかを確認しました。 特に効いたのは easily_misread_points です。二軸グラフ、単位が%の軸、近接したデータラベル。「ここを外したら読めていないことになる」という点を先に洗い出しておくと、3方式の出力を並べたときに何を見ればよいかが決まります。 結果 3方式の総合比較 3試行の平均値です。 方式 リクエスト数 1枚あたり facts件数 処理時間 入力 トークン 出力 トークン コスト ① PDF一括 2 2.14 195秒 44,896 19,619 $0.29 ② 3枚ずつ 9 4.33 573秒 128,980 48,609 $0.74 ③ 1枚ずつ 21 5.98 816秒 185,388 62,901 $1.00 この結果から、以下が分かります。 1枚あたりに読み取る事実(facts)の件数は ③ > ② > ① の順にきれいに並んだ 。刻みを細かくするほど細かく読む、という素朴な予想通りの結果 ①のfacts件数は明確に少ない 。2.14件は②の約半分、③の約1/3 ②と③の差はほぼ比例 。②→③は件数1.38倍に対し、時間1.4倍・コスト1.35倍。細かく刻んだ分だけコストを払っている 処理速度とコストは ① > ② > ③ 。①→③で時間4.2倍・コスト3.5倍 時間・コストの内訳と、比較上の注意(クリックで展開) 上表の内訳は以下です(1試行あたり、API実測値)。 方式 解析ステップ 最終統合 ① PDF一括 1回 133秒 / $0.197 1回 55秒 / $0.089 ② 3枚ずつ 7回 326秒 / $0.471 2回 165秒 / $0.273 ③ 1枚ずつ 19回 520秒 / $0.623 2回 195秒 / $0.377 ②③の統合が2回になっているのは、検証途中で統合ステップの出力上限不足(截断)が見つかり、公平性のため全試行を上限を上げて再実行したためです。失敗した旧リクエストもログに残しているので、②③の計測値には2回分が含まれています(①は上限修正後に実行したため1回)。 統合1回分に揃えて概算すると、時間は195秒 / 409秒 / 617秒、コストは$0.29 / $0.61 / $0.81となり、①→③の倍率は 時間3.3倍・コスト2.8倍 になります。倍率の大小は変わりますが、順序と桁感は変わりません。 処理時間は壁時計ベースです。検証中にPCがスリープした試行があったため、API計測時間との乖離を検知して除外した頑健値を使っています。 料金は検証時点のSonnet 5の導入価格(入力$2.00 / 出力$10.00 per 1M tokens)で計算しています。 増えた分に何が入っていたか 件数の話はここまでです。ここからは、その差が中身を伴っているのかを、実際のスライドと出力内容で見ていきます。 スライド9 まずは、スライド9です。 このスライドは、赤線の総人口だけが右軸で、他の値は左軸で表現されています。 各手法でのimportant_numbers(重要指標)の結果は以下です。 ①<②<③と、抽出できた重要指標の数が増えていることが分かります。 特に、①では左右の軸に関する言及はありませんが、②と③ではしっかりと言及がされています。 種別 結果 ① PDF一括 [総人口推計:1995年度〜2045年(実績値・推計値、薩摩川内市)] ② 3枚ずつ ["対象期間:1995年度〜2045年(実績値は〜2020年頃まで、以降は推計値)","総人口軸:右軸、目盛り約65,000〜115,000人","出生・死亡・転入・転出数:左軸、目盛り約1,000〜6,500人"] ③ 1枚ずつ ["対象期間:1995年度頃〜2045年(実績値と推計値を含む、薩摩川内市)","総人口軸(右軸):目盛65,000人〜115,000人の範囲で表示(薩摩川内市、期間中の推移)","出生数・死亡数・転入数・転出数軸(左軸):目盛500人〜6,500人の範囲で表示(薩摩川内市、期間中の推移)","2010年度以降、自然減(死亡数が出生数を上回る状態)が拡大傾向(薩摩川内市)","近年、転入・転出ともに増加傾向だが一定の転出超過が継続(薩摩川内市)"] スライド7 次は、スライド7です。 このスライドは、縦軸の単位が%であり、実数ではないことが注意点です。 各手法でのconditions_and_notes(条件や備考)の結果は以下です。 こちらもスライド9同様に①<②<③と、抽出した備考の数が増えていることが分かります。 縦軸の単位が%であることに注意を促しているのは、③のみです。 種別 結果 ① PDF一括 ["年少人口=0〜14歳、生産年齢人口=15〜64歳、老年人口=65歳以上(基礎知識欄)。"] ② 3枚ずつ ["グラフは1985年〜2020年頃までが実績値、それ以降が推計値である(グラフ下部の矢印表示)。", "年少人口:0〜14歳、生産年齢人口:15〜64歳、老年人口:65歳以上と定義されている。", "本スライドは自然増減(出生・死亡)と社会増減(転入・転出)の内訳分解には踏み込んでおらず、年齢3区分の増減率のみを扱っている。"] ③ 1枚ずつ ["グラフの縦軸は増減率(%)であり、実数(人数)ではない点に留意が必要。", "2020年以前が実績値、2020年以降(2025年〜2045年)が推計値であることがグラフ下部の矢印で明示されている。", "グラフ上の具体的な数値ラベルは示されておらず、目視での概算読み取りにとどまるため、正確な増減率数値は資料からは確定できない。", "本スライドは特定自治体(薩摩川内市)を用いた「分析例」であり、教材としての位置づけであることがヘッダーのナビゲーション(基礎分析:人口構成→人口増減→自然増減→社会増減→将来人口推計、応用分析:人口構成背景→出生数の増減要因→転入転出要因)から読み取れる。"] これらの結果を見ると、① PDF一括は抽出内容が大雑把になりがちで、一方、②と③は詳細を読み解こうとしていることが分かります。 さらに②よりも③の方が、細かく指標や条件を抽出しようとしていることが分かります。 つまり、facts件数の差は水増しではなく、 軸・単位・注記という「外すと読み違える」項目の差 でした。ゴールドの easily_misread_points に挙げた罠のうち、二軸グラフの識別は②から拾われ、軸の単位まで届くのは③のみ、という並びになっています。 使い分けの指針 結局のところ、 どこまで細部が必要かで選べばよい という話になります。 ① PDF一括 :大雑把に、早く内容を理解したいとき ② 3枚ずつ :詳細に把握したいが、完全な細部までは不要なとき ③ 1枚ずつ :時間をかけても、詳細までしっかり把握したいとき そのうえで、入れておいた方がいい実装上の備えを挙げておきます。 スライド単位の中間出力を必ず持つ 。PDFから直接まとめさせる構成も試しましたが、根拠として引用するスライド番号がズレました(対象外のスライド1〜4を引用する誤りが3試行中2試行で発生)。解析ステップを1段挟むだけで、根拠スライド番号の妥当性は0.886→1.000になりました スライド番号の欠落を機械的に検知する 。まとめて渡す方式は、1リクエストの失敗がブロック全体の欠落になります 出力上限は実測で決める 。当初は解析2,000 / 統合3,000で設定していましたが、実測では1スライドの解析ですら截断されました。最終的に解析6,000 / 統合20,000(PDF一括の解析のみ32,000)です クライアントのタイムアウトを確認する 。PDF一括の解析は128〜136秒かかり、SDKの既定タイムアウト120秒に足りず3回連続で失敗しました。長い入力を1回で投げる方式では、まず疑うべきポイントです 不確実性の明示を許すプロンプトにする 。「読み取れない箇所は推測せず unclear_points に書く」と指示しておくと、スライド11のような罠で断定を回避してくれます おまけ:PPTXをPNG化すると画像がズレる 最後に、LLMに渡す前段でハマった話です。 画像化の処理をLLMに書かせたら、当然のようにLibreOffice( soffice --headless --convert-to pdf → pdftoppm )を使う実装が出てきました。ところがこれ、 資料のフォントが環境に無いとフォント代替が起きてレイアウトが崩れます 。今回の資料は Meiryo UI などを使っていて未埋め込みだったため、本文の行数が増え、一部スライドで下の見出し帯にテキストがはみ出して重なりました。入力画像がズレていては、比較検証にすらなりません。 回避するには、 フォントを持っているアプリに描画させる のが確実です。macOSの Microsoft PowerPoint.app はアプリ内部( Contents/Resources/DFonts/meiryo.ttc )にMeiryoを同梱しているので、ここを通せば崩れません。 手作業でよければ :PowerPointから直接PNGエクスポート 自動化するなら :PowerPointでPDFに書き出してから、 pdftoppm で1枚ずつPNG化 今回は後者を使いました。資料側でフォントを埋め込んでもらえるなら、それが最善です。 おわりに 今回は、情報量の多い日本語行政資料をClaude Sonnet 5に読ませるとき、入力の刻み方だけで読み取りの細かさ・時間・コストがどう変わるかを実測しました。 資料全体を統合したまとめは、どの刻み方でもほぼ同じ。差が出るのは1ページごとの内容理解の細部 その細部の読み取り量は「1枚ずつ > 3枚ずつ > PDF一括」の順で、1枚あたり5.98件 / 4.33件 / 2.14件 増えた分は水増しではなく、軸・単位・注記という「外すと読み違える」項目だった 1枚ずつは最も細かく読むが、PDF一括に比べて時間4.2倍・コスト3.5倍 PDF一括が落とすのは数値そのものより、注記・例外条件・軸の但し書き。しかもこれは出力上限の問題ではなく、モデルが自ら記述を圧縮する挙動 3枚ずつは「まとめて1回」をやめる効果の半分以上を取れる。ただしブロック欠落のリスク管理が必要 図の情報量そのものに起因する誤読は、刻み方を変えても消えない 個人的にいちばん意外だったのは、 方式ごとに強みの質が違った ことです。 細かく刻むと数値の網羅性が上がり、まとめて渡すと資料全体の枠組みを取り違えにくくなる。「精度」という一つの数字に丸めると、この非対称性が見えなくなってしまいます。実運用では、一括入力で全体の枠組みを掴んだうえで、数値の根拠が必要なスライドだけ1枚ずつ精読させる、という組み合わせも十分にありだと思います。 もちろん限界もあります。対象は1つの資料(19スライド)だけで、各方式3試行のみ。中身の確認は筆者の目視です。資料の種類(文字中心の議事録、表中心の財務資料など)が変われば、最適な刻み方も変わるはずです。 それでも、「1ページずつか、まとめてか」を感覚ではなく数字で決めたいときに、当たりをつける材料にはなるかと思います。本記事が、誰かの参考になれば幸いです。 出典・注意事項 出典 出典:RESAS Portal - 地域課題分析ナビゲーション テーマ① 地域の人口減少対策 -( https://resas-portal.go.jp/region/ ) 加工内容 本記事では、上記資料のスライド5〜23を検証対象とし、PNG画像化・PDF結合・縮小を行った上でLLMへ入力しています。 掲載しているスライド画像は、検証目的での画像化・縮小を行ったものです(RESAS Portalを加工して作成)。 記事化にあたっては、サイトポリシー(政府標準利用規約準拠)を確認のうえ、問い合わせ窓口へメールで確認し、利用の許諾をいただいています。 免責・注意事項 本記事は筆者による技術検証の内容であり、当社での正式な推奨手法や導入事例を示すものではありません。実際の業務適用にあたっては、各組織のガイドラインに従ってください。 本検証はLLMの内容理解支援力を確認する技術検証であり、対象資料や特定自治体の政策を評価・決定するものではありません。 LLMの出力は筆者による技術検証の結果であり、内閣府・RESAS Portal・資料内で言及される自治体の見解ではありません。 資料内に掲載された分析例を、対象地域の最新の人口動態や政策状況として扱わないでください。 人口移動や出生に関する分析結果を、特定の年代・性別・個人の責任として解釈しないでください。 料金は検証時点(2026年8月)のClaude Sonnet 5の価格に基づく実測値です。最新の価格は公式ドキュメントをご確認ください。
はじめに こんにちは、サイオステクノロジーの小沼 俊治です。 「理屈はいいから、まずは実際に Apache Kafka やストリーム処理というものを動かして体験してみたい」。 本記事は、そんな方々に向けて、リアルタイムデータ基盤に不可欠なイベントストリーミングやデータパイプラインの基礎を手を動かしながら学習できる、実践的な入門ガイドとして用意しました。 単なるメッセージの送受信にとどまらず、リアルタイムなデータ加工(Flink / Kafka Streams)、データの品質管理(Schema Registry)、データベースとの自動同期(CDC Connector)、そしてパーティションによる負荷分散までを包括したリアルタイムデータパイプライン環境を無料で体験できるハンズオンを提供します。 本ハンズオンでは、以下のオープンソース・プロダクトのみで構成された環境をコンテナを使って一括で立ち上げ、データが発生してから加工・連携・分散処理されるまでの一連の流れを体験していただきます。 Apache Kafka:イベントストリーミング基盤(Broker) AKHQ:Kafka クラスタ管理 Web UI Apache Flink:分散ストリーム処理エンジン(Flink SQL / Table API / DataStream API) Apicurio Registry:スキーマ管理・データ品質保持(Schema Registry) Kafka Streams:組み込み型ストリーム処理アプリケーション Debezium Connector for MySQL:データベース変更データ捕獲(CDC) MySQL:データベース なお、本記事は「Kafka を中心としたエコシステムによるリアルタイム処理の流れ」を体験いただくことを主目的としています。そのため、環境構築手順そのものの詳細な解説は一部割愛しており、「まずは動かしてみたい」という方に最適です。もし環境構築の裏側に興味を持っていただいた場合は、後半の Appendix や GitHub リポジトリの設定ファイルをぜひ解析してみてください。 構成概要 筆者が動かした際の主な構成要素は以下の通りです。 Windows 11 Professional WSL 2.4.12.0 Ubuntu 24.04.2 LTS Docker Engine 28.0.4 Apache Kafka 4.3.1 Windows (WSL) 以外のOSをご利用の方も、条件が満たしていれば以下の手順からハンズオンを進められます。 Ubuntu (Linux) 環境の方: WSL の構築は不要なため、「 Docker Engine 環境の構築 」章から開始してください。 macOS 環境の方: Docker Desktop for Mac などでコンテナ実行環境が準備済みであれば、「 ハンズオンに必要なコマンドの準備 」章から開始してください。 ハンズオンを構成する環境は以下の通りです。 Kafka のメインコンポーネントを担う Broker と、動作状況や設定を Web UI で提供する AKHQ をコンテナで構築します。 ストリーム処理を体験するため以下のコンテナを構築します。 Apache Flink の Flink SQL、Table API、および DataStream API を体験するため、jobmanager、taskmanager をコンテナで構築します。 Streams Application を体験するため、Java 環境のコンテナを構築します。 Schema Registry を体験するため、スキーマを管理する Apicurio Registry をコンテナで構築します。 Kafka Connector を体験するため、データベースから更新差分を転送する CDC をコンテナで構築します。 データベースとして MySQL をコンテナで構築します。 CDC Connector の稼働環境として Debezium connector をコンテナで構築します。 トピックでメッセージを送受信するために、ローカル環境に Java や Kafka CLI ツールをインストールする手間を省くため、プロデューサーとコンシューマーを実行する専用の「handson-client」コンテナを用意しています。 環境構築や各種設定に使用するそれぞれのファイルは、以下の GitHub リポジトリで公開しています。 https://github.com/Toshiharu-Konuma-sti/hands-on-kafka $ tree ~/handson/hands-on-kafka/ hands-on-kafka/ |-- container/ …… 「環境構築」章でコンテナ作成で使う素材 | |-- docker-compose.yml | : | |-- development/ | `-- streams/ …… 「独自 Streams アプリケーション開発」章で使うアプリケーションの開発素材 | |-- setup/ …… 「kafka 演習」章などのハンズオン環境の構築に必要な各種設定素材 | |-- SETUP_HANDS-ON.sh | : | `-- try-my-hand/ …… 「Kafka 演習」章などのハンズオンで使う環境 : 基礎環境の構築 WSL 環境の構築 Windows PC の場合には、 以下手順を参考に WSL と Linux ディストリビューション(Ubuntu)環境を用意します。 初期環境構築: WSL 環境 on Windows Docker Engine 環境の構築 コンテナ環境を使うため、以下手順を参考に Ubuntu へ Docker Engine 環境を用意します。 初期環境構築: Docker Engine on Ubuntu ハンズオンに必要なコマンドの準備 本ハンズオンの実施には、以下のコマンドやランタイムが必要です。 これらは主に、環境構築を行う際に使用します。 jq コマンド:API へ投入するリクエストの JSON 構造の編集をしたり、JSON 形式のレスポンスから特定の値を抽出・整形するために使用します。 利用箇所: Kafka 環境セットアップスクリプトの実行 インストールされていない場合は、以下手順を参照して Ubuntu 環境へ用意します。 初期環境構築: ユーティリティツール on Ubuntu Kafka 環境の構築 GitHub からハンズオン用のリポジトリ取得 ハンズオンを進めるための環境構築用の設定ファイルやスクリプトを含んだリポジトリを GitHub からダウンロードして取得します。 GitHub からハンズオン用のリポジトリ取得 $ mkdir -p ~/handson/ $ cd ~/handson/ 「 $ git clone 」コマンドで本ハンズオン用のリポジトリを取得します。 $ git clone https://github.com/Toshiharu-Konuma-sti/hands-on-kafka.git $ cd hands-on-kafka/ コンテナ構築スクリプトの実行 本章ではターミナルを用いて以下のディレクトリで作業を実施します。 $ cd ~/handson/hands-on-kafka/container/ コンテナ構築用に用意してあるスクリプトを実行して、Kafka 環境の各種コンテナを構築します。 $ ./CREATE_CONTAINER.sh info /************************************************************ * Information: * - Navigate to Web ui tools with the URL below. * - AKHQ: http://localhost:9021 * - Flink dashboard: http://localhost:8181 ***********************************************************/ 「AKHQ」コンテナが稼働してブラウザでアクセスできます。 http://localhost:8080 なお、コンテナ構築スクリプトで実行する内容は以下を参照してください。 コンテナ構築スクリプトの解説 Kafka 環境セットアップスクリプトの実行 トピックの作成やスキーマ登録などを、スクリプトを使って一気に行います。これにより、複雑な設定を手動で行う手間を省きます。 本章ではターミナルを用いて以下のディレクトリで作業を実施します。 $ cd ~/handson/hands-on-kafka/setup/ セットアップ用に用意してあるスクリプトを実行して、Kafka 環境の各種設定を行います。 $ ./SETUP_HANDS-ON.sh コマンドが足りずにスクリプトが終了した際は、以下を参照してインストールしてください。 初期環境構築: ユーティリティツール 設定スクリプトで実行する内容は以下を参照してください。 Kafka 環境セットアップスクリプトの解説 Kafka の演習 ここから始まる Kafka のハンズオンは「 try-my-hand/ 」ディレクトリで実施します。 $ cd ~/handson/hands-on-kafka/try-my-hand/ なお、プロデューサーやコンシューマーの実行に必要な Java や Python 環境は、あらかじめ「handson-client」コンテナ内に用意されています。手元のローカル環境に各種ランタイムをインストールする必要はありません。 今後のハンズオンに関するコマンド操作は、すべて「handson-client」コンテナ経由で実行します。 コマンド例) $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --create --topic sample-topic トピック操作 メッセージの送受信に必要なトピックの作成や削除を始めとする操作を体験します。 トピックを作成します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --create --topic sample-topic Created topic sample-topic. 存在するトピック一覧を確認します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --list : sample-topic トピックを削除します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --delete --topic sample-topic 再度トピックを作成して、パーティションを分割します(例は3分割)。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --create --topic sample-topic $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --alter --topic sample-topic --partitions 3 存在するトピックの詳細を確認します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --describe --topic sample-topic Topic: sample-topic TopicId: M5GoT1uKRPud0AFsZTAdWg PartitionCount: 3 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: sample-topic Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: Topic: sample-topic Partition: 1 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: Topic: sample-topic Partition: 2 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: トピックを削除して次に進みます。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --delete --topic sample-topic メッセージ送受信 トピックを介してプロデューサーからコンシューマーへメッセージの送受信を体験します。 2つのターミナルを立ち上げて、プロデューサー(送信側)とコンシューマー(受信側)間でメッセージの送受信を体験します。 送信側のターミナルでプロデューサーを起動してメッセージの送信待機状態にし、次に受信側のターミナルでコンシューマーを起動してメッセージの受信待機状態にします。それぞれが待機状態になっている状況で、送信側のターミナルからメッセージをタイピングして送信した後に、受信側に届いたことを確認して進めます。 まずは、「 Kafka 環境セットアップスクリプトの実行 」章で作成した、本章で利用するトピックが存在することを確認します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --describe --topic my-topic Topic: my-topic TopicId: dlvh0VszQCi0we8mhEwnwg PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-topic Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: 基本的な送受信 送信側のターミナルで送信先のトピックを指定してプロデューサーを起動すると、メッセージ送信の待機状態になります。なお、待機状態は「Ctrl + C」で解除できます。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-topic > 受信側のターミナルで受信元のトピックを指定してコンシューマーを起動すると、メッセージ受信の待機状態になります。なお、待機状態は「Ctrl + C」で解除できます。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-topic 送信側のプロデューサーからメッセージを、1行ごとに1メッセージを送信します。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-topic > hello001 受信側のコンシューマーにメッセージが届きます。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-topic hello001 受信側で「Ctrl + C」を押してコンシューマーを終了して、再度コンシューマーを起動します。 【受信側】 [Ctrl + C] $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-topic 送信側のプロデューサーから追加でメッセージを送信します。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-topic > hello001 > hello002 受信側のコンシューマーでは過去のメッセージは得られず、受信状態中に送信されたメッセージのみが受信できます。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-topic hello002 ここまでで送受信したメッセージが、AKHQ からも確認できます。 http://localhost:8080/ui/hands-on-kafka/topic/my-topic/data 最初のメッセージから受信 受信側で「Ctrl + C」して、「–from-beginning」オプションを付与してコンシューマーを起動すると、トピック内の最初のメッセージから受信できます。 【受信側】 [Ctrl + C] $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-topic \ --from-beginning hello001 hello002 送信側のプロデューサーから追加でメッセージを送信します。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-topic > hello001 > hello002 > hello003 受信側のコンシューマーに追加のメッセージが届きました。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-topic \ --from-beginning hello001 hello002 hello003 メッセージ情報を付与して受信 受信側で「Ctrl + C」して、メッセージに関する各種情報を表示するオプションを付与してコンシューマーを起動します。メッセージの送信時間やオフセット(順番)などが表示されます。 【受信側】 [Ctrl + C] $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-topic \ --formatter-property print.timestamp=true \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.headers=true \ --formatter-property print.key=true \ --from-beginning CreateTime:1786604356842 Partition:0 Offset:0 NO_HEADERS null hello001 CreateTime:1786613536930 Partition:0 Offset:1 NO_HEADERS null hello002 CreateTime:1786622414755 Partition:0 Offset:2 NO_HEADERS null hello003 オフセット指定でメッセージ受信 受信側で「Ctrl + C」して、オフセットを指定してコンシューマーを起動することで、指定位置からメッセージを受信できます。 【受信側】 [Ctrl + C] $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-topic \ --formatter-property print.timestamp=true \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.headers=true \ --formatter-property print.key=true \ --partition 0 --offset 1 CreateTime:1786613536930 Partition:0 Offset:1 NO_HEADERS null hello002 CreateTime:1786622414755 Partition:0 Offset:2 NO_HEADERS null hello003 Key & Value 形式で送受信 今度は送信側で「Ctrl + C」して、Key & Value 形式でメッセージを送信するオプションを付与してコンシューマを起動します。例では、「:」をデリミタにして Key & Value 形式でメッセージを送信しています。 【送信側】 [Ctrl + C] $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-topic \ --reader-property "parse.key=true" --reader-property "key.separator=:" > key1:value1 受信側のコンシューマーに Key & Value 形式で追加のメッセージが届きました。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-topic \ --formatter-property print.timestamp=true \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.headers=true \ --formatter-property print.key=true \ --partition 0 --offset 1 CreateTime:1786613536930 Partition:0 Offset:1 NO_HEADERS null hello002 CreateTime:1786622414755 Partition:0 Offset:2 NO_HEADERS null hello003 CreateTime:1786626086546 Partition:0 Offset:3 NO_HEADERS key1 value1 Apache Flink 本章では、Kafka に流れるデータをリアルタイムに加工するストリーム処理エンジン「Apache Flink」を体験します。 リアルタイムデータ変換においてストリーム処理を行う手法として、「Flink SQL」「Table API」「DataStream API」の3つのアプローチを順に試していきます。 Flink SQL プロデューサーからコンシューマーへメッセージを送受信する過程で、事前に DDL で登録したテーブルと SQL クエリを用いて、リアルタイムにデータをストリーム処理する Flink SQL を体験します。 「my-flink-input」トピックから受信して「my-flink-sql-output」トピックへ送信する過程で、Flink では Flink SQL を使って以下の加工を行います。 「my-flink-input」トピックから流れてきた JSON 形式のメッセージをフィールド分解し、「FLINK_SQL_INPUT」テーブルのレコードとして保存します。 「FLINK_SQL_INPUT」テーブルのレコードを「FLINK_SQL_OUTPUT」テーブルに挿入します。挿入の際には以下の加工を行います。 対象フィールドを「name」と「gender」に絞ります。 「gender」フィールドの値が 'M' または 'X' のデータのみに絞ります。 本ハンズオンでストリーム処理として実装した SQL: flink_sql.sql まずは、「 Kafka 環境セットアップスクリプトの実行 」章で作成した、本章で利用するトピックが存在することを確認します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --describe \ --topic my-flink-input,my-flink-sql-output Topic: my-flink-sql-output TopicId: YIyil35cQe2TK-3b0yTPgA PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-flink-sql-output Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: Topic: my-flink-input TopicId: GAKyq9SZSNmV9_LRP60qug PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-flink-input Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: メッセージ送受信 送信側のターミナルで、Flink SQL を体験するための送信先のトピックを指定してプロデューサーを起動します。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-flink-input > 受信側のターミナルで、Flink SQL を体験するための受信元のトピックを指定してコンシューマーを起動します。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-flink-sql-output \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.key=true \ --from-beginning 送信側のプロデューサーから JSON 形式のメッセージ送信します。(Flink SQL で加工前) 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-flink-input > {"id":1, "name":"taro", "gender":"M", "age":10} > {"id":2, "name":"hanako", "gender":"F", "age":20} > {"id":3, "name":"tama", "gender":"X", "age":30} > {"id":4, "name":"tamako", "gender":"F", "age":40} 受信側のコンシューマーに、Flink SQL で加工後のメッセージが届きました。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-flink-sql-output \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.key=true \ --from-beginning Partition:0 Offset:0 null {"name":"taro","gender":"M"} Partition:0 Offset:1 null {"name":"tama","gender":"X"} Flink SQL で実行された ETL 処理の結果を確認します。 Flink SQL のカタログ情報はセッションごとのインメモリ管理となるため、起動した SQL Client 上で一度 CREATE TABLE 文を実行してから SELECT を発行します。 $ docker exec -it jobmanager ./bin/sql-client.sh Flink SQL> SHOW TABLES; Empty set Flink SQL> CREATE TABLE flink_sql_input ( id INT, name STRING, gender STRING, age INT ) WITH ( 'connector' = 'kafka', 'topic' = 'my-flink-input', 'properties.bootstrap.servers' = 'kafka:29092', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' ); [INFO] Execute statement succeeded. Flink SQL> CREATE TABLE flink_sql_output ( name STRING, gender STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'my-flink-sql-output', 'properties.bootstrap.servers' = 'kafka:29092', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' ); [INFO] Execute statement succeeded. Flink SQL> SHOW TABLES; +------------------+ | table name | +------------------+ | flink_sql_input | | flink_sql_output | +------------------+ 2 rows in set Flink SQL> SELECT * FROM flink_sql_input; id name gender age --- ------ ------ --- 1 taro M 10 2 hanako F 20 3 tama X 30 4 tamako F 40 Flink SQL> SELECT * FROM flink_sql_output; name gender ------ ------ taro M tama X Flink SQL> quit; プロデューサーの起動、コンシューマーの起動、および Flink SQL で ETL 処理の確認は、簡単に実行できるように以下スクリプトを用意しています。 プロデューサーの起動: flink-producer.sh コンシューマーの起動: flink-sql-consumer.sh Flink SQL で ETL 処理の確認: flink-sql-check-tables.sh Table API (PyFlink) プロデューサーからコンシューマーへメッセージを送受信する過程で、事前に DDL でスキーマを登録する代わりに、プログラム内で動的にスキーマを定義・操作しながら、柔軟かつリアルタイムにデータをストリーム処理する Table API を体験します。 なお、本章の Table API のプログラムは、Python 言語版である PyFlink を利用して実装しています。 「my-flink-input」トピックから受信して「my-flink-table-api-output」トピックへ送信する過程で、Flink では Table API を使って以下の加工を行います。 「my-flink-input」トピックから流れてきた JSON 形式のメッセージをフィールド分解し、「FLINK_TABLE_API_INPUT」テーブルのレコードとして保存します。 「FLINK_TABLE_API_INPUT」テーブルのレコードを「FLINK_TABLE_API_OUTPUT」テーブルに挿入します。挿入の際には以下の加工を行います。 「name」フィールドの大文字と小文字を反転(相互変換)します。 「name」フィールドに、「gender」フィールドが 'M' の場合は 'Mr.' 、 'F' の場合は 'Ms.' 、および 'X' の場合は後方に '-san' の敬称を付与します。 「gender」フィールドを大文字を小文字に、小文字を大文字に変換します。 本ハンズオンでストリーム処理として実装した Table API プログラム: table_api.py まずは、「 Kafka 環境セットアップスクリプトの実行 」章で作成した、本章で利用するトピックが存在することを確認します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --describe \ --topic my-flink-input,my-flink-table-api-output Topic: my-flink-input TopicId: zqLyeMLwQ6aDptavQVkrvg PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-flink-input Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: Topic: my-flink-table-api-output TopicId: E7t1BHUSRQKYAaNQF8Fp6A PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-flink-table-api-output Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: メッセージ送受信 送信側のターミナルで、Flink Table API を体験するための送信先のトピックを指定してプロデューサーを起動します。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-flink-input > 受信側のターミナルで、Flink Table API を体験するための受信元のトピックを指定してコンシューマーを起動します。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-flink-table-api-output \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.key=true \ --from-beginning 送信側のプロデューサーから JSON 形式のメッセージ送信します。(Flink Table API で加工前) 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-flink-input > {"id":5, "name":"jiro", "gender":"M", "age":50} > {"id":6, "name":"saki", "gender":"F", "age":60} > {"id":7, "name":"pochi", "gender":"X", "age":70} 受信側のコンシューマーに、Flink Table API で加工後のメッセージが届きました。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-flink-table-api-output \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.key=true \ --from-beginning Partition:0 Offset:4 null {"id":5,"name":"Mr. JIRO","gender":"m","age":50} Partition:0 Offset:5 null {"id":6,"name":"Ms. SAKI","gender":"f","age":60} Partition:0 Offset:6 null {"id":7,"name":"POCHI -san","gender":"x","age":70} Flink Table API で実行された ETL 処理の結果を確認します。 Flink Table API のカタログ情報はセッションごとのインメモリ管理となるため、起動した SQL Client 上で一度 CREATE TABLE 文を実行してから SELECT を発行します。 $ docker exec -it jobmanager ./bin/sql-client.sh Flink SQL> SHOW TABLES; Empty set Flink SQL> CREATE TABLE flink_table_api_input ( id INT, name STRING, gender STRING, age INT ) WITH ( 'connector' = 'kafka', 'topic' = 'my-flink-input', 'properties.bootstrap.servers' = 'kafka:29092', 'properties.group.id' = 'pyflink-group', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' ); [INFO] Execute statement succeeded. Flink SQL> CREATE TABLE flink_table_api_output ( id INT, name STRING, gender STRING, age INT ) WITH ( 'connector' = 'kafka', 'topic' = 'my-flink-table-api-output', 'properties.bootstrap.servers' = 'kafka:29092', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' ); [INFO] Execute statement succeeded. Flink SQL> SHOW TABLES; +------------------------+ | table name | +------------------------+ | flink_table_api_input | | flink_table_api_output | +------------------------+ 2 rows in set Flink SQL> SELECT * FROM flink_table_api_input; id name gender age --- ------ ------ --- 5 jiro M 50 6 saki F 60 7 pochi X 70 Flink SQL> SELECT * FROM flink_table_api_output; id name gender age --- ----------- ------ --- 5 Mr. JIRO M 50 6 Ms. SAKI F 60 7 POCHI -san X 70 Flink SQL> quit; プロデューサーの起動、コンシューマーの起動、および Flink SQL で ETL 処理の確認は、簡単に実行できるように以下スクリプトを用意しています。 プロデューサーの起動: flink-producer.sh コンシューマーの起動: flink-table-api-consumer.sh Table API で ETL 処理の確認: flink-table-api-check-tables.sh DataStream API (PyFlink) プロデューサーからコンシューマーへメッセージを送受信する過程で、流れてくるデータを表形式ではなく1件ずつの「オブジェクト」として直接操作し、状態管理(Stateful)や柔軟なタイムウィンドウ処理など、より細やかで高度な制御を行いながらストリーム処理する DataStream API を体験します。 なお、本章の DataStream API のプログラムは、Python 言語版である PyFlink を利用して実装しています。 「my-flink-input」トピックから受信して「my-flink-datastream-api-output」トピックへ送信する過程で、Flink では DataStream API を使って以下の加工を行います。 「my-flink-input」トピックから流れてきた JSON 形式のメッセージを1件ずつ、以下の加工を順次処理します。 「name」と「gender」フィールドへ Table API と同じ加工を行います。 Table API では実装が困難な加工処理も行います。 ストリーム処理時点における性別毎の平均年齢値を保持し、”avg_age” フィールドに出力します。 性別毎の平均年齢値とストリームのレコードを比較した結果を “vs_avg“ フィールドに出力します。 加工処理が終わったら、「my-flink-datastream-api-output」トピックと紐づいた「FLINK_DATASTREAM_API_SINK」テーブルにシンクします。 本ハンズオンでストリーム処理として実装した DataStream API プログラム: datastream_api.py まずは、「 Kafka 環境セットアップスクリプトの実行 」章で作成した、本章で利用するトピックが存在することを確認します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --describe \ --topic my-flink-input,my-flink-datastream-api-output Topic: my-flink-datastream-api-output TopicId: wYDXtcFLR6GZf_qJ53POTA PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-flink-datastream-api-output Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: Topic: my-flink-input TopicId: rAlyReroQFyQprZbtUJUOA PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-flink-input Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: メッセージ送受信 送信側のターミナルで、Flink DataStream API を体験するための送信先のトピックを指定してプロデューサーを起動します。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-flink-input > 受信側のターミナルで、Flink DataStream API を体験するための受信元のトピックを指定してコンシューマーを起動します。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-flink-datastream-api-output \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.key=true \ --from-beginning 送信側のプロデューサーから JSON 形式のメッセージ送信します。(Flink DataStream API で加工前) 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-flink-input > {"id":8, "name":"saburo", "gender":"M", "age":80} > {"id":9, "name":"yoko", "gender":"F", "age":90} > {"id":10, "name":"kuro", "gender":"X", "age":100} 受信側のコンシューマーに、Flink DataStream API で加工後のメッセージが届きました。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-flink-datastream-api-output \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.key=true \ --from-beginning Partition:0 Offset:7 null {"id":8,"name":"Mr. SABURO","gender":"m","age":80} Partition:0 Offset:8 null {"id":9,"name":"Ms. YOKO","gender":"f","age":90} Partition:0 Offset:9 null {"id":10,"name":"KURO -san","gender":"x","age":100} プロデューサーの起動、コンシューマーの起動、および Flink SQL で ETL 処理の確認は、簡単に実行できるように以下スクリプトを用意しています。 プロデューサーの起動: flink-producer.sh コンシューマーの起動: flink-datastream-api-consumer.sh DataStream API で扱うシンクテーブルの確認: flink-datastream-api-check-table.sh Flink 各種ストリームの使い分け方 ストリーム処理別の相違点 ストリーム処理方式 メッセージの処理方法 処理の実装媒体 処理の組み込み方法 ストリーム処理方式 トピックをテーブルとして処理 SQL文(DDL、DML) Flink SQL Client で SQL を実行 Table API トピックをテーブルとして処理 ソースコード jobmanager でソースコードを実行 DataStream API 1件ずつのメッセージで処理 ソースコード jobmanager でソースコードを実行 Table API と DataStream API の相違点 ソースコードで実装する Table API と DataStream API の違いは以下の通りです。 入力元との接続方式: Table API: TableEnvironment.create() で生成したコンテキストに対して、 execute_sql() の DDL で Kafka トピックにスキーマ(型定義)を割り当て、仮想テーブルとして登録してから処理を開始します。 DataStream API: StreamExecutionEnvironment.get_execution_environment() で生成したコンテキストに対して、 KafkaSource.builder() で構築したコネクションを渡し、ストリームとして直接取り込みます。入力に DDL は使いません。 データ処理の書き方と単位: Table API:データを表として扱うため、 @udf デコレータを付けた関数を定義し、 .select() や .filter() などのメソッドチェーンで列単位の変換や抽出を行います。 DataStream API:データを「1件のメッセージ(オブジェクト)」として受け取ります。シンプルな変換では map() や flatMap() などのオペレータを用いて変換・抽出を行いますが、複雑な条件判断や状態管理を行う場合は KeyedProcessFunction や process() を実装することで、キーごとに状態(State)を保持したより高度で柔軟なロジックを記述できます。 出力先との接続方法: Table API:入力と同様に DDL で定義した出力先テーブルに対し、 execute_insert() を発行してデータを流し込みます。 DataStream API: KafkaSink などの専用コネクタを定義し、 .sink_to() メソッドでストリームを直接出力先へ接続します。なお、出力部のみ Table API の DDL 機能をブリッジとして利用するハイブリッドなパターンも存在します。 本ハンズオンで動かしているプログラムは以下の通りです。 Table API: table_api.py DataStream API: datastream_api.py Schema Registry プロデューサーからコンシューマーへメッセージを送受信する過程で、一元管理されたデータ構造(スキーマ)とメッセージを照合し、ルールに合わない不正データの混入を防いでシステム全体のデータ品質を守る Schema Registry を体験します。 プロデューサーからコンシューマーへトピックを介してメッセージを送受信する過程で、以下の照合処理を行います。 プロデューサーでメッセージを送信する際に以下の処理を行います。 事前に Schema Registry にスキーマ登録した際に発行された ID を指定して、Schema Registry に登録されているスキーマを取得します。 取得したスキーマとメッセージを照合し、一致したメッセージのみスキーマ ID を付与してトピックへ送信します。 コンシューマーでメッセージを受信する際に以下の処理を行います。 トピックを介して送られてきたメッセージに付与された ID を指定して、Schema Registry に登録されているスキーマを取得します。 取得したスキーマとメッセージを照合し、一致したメッセージのみ受取ります。 Kafka 自体や、Kafka CLI ツールのプロデューサーやコンシューマーには、スキーマとメッセージを照合する仕組みは備わっていないため、本章では独自に実装したプロデューサーとコンシューマーを使用しています。 schema-registry-json-producer.py 、 schema-registry-json-consumer.py Schema Registry として利用している Apicurio Registry(kafkasql ストレージモード) では、登録されたスキーマを Kafka の「kafkasql-journal」トピックに保存して管理します。 本ハンズオンで Schema Registry 登録用に実装したスキーマ: schema_registry.json AKHQ で登録されているスキーマを確認できます。 http://localhost:8080/ui/hands-on-kafka/schema まずは、「 Kafka 環境セットアップスクリプトの実行 」章で作成した、本章で利用するトピックが存在することを確認します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --describe \ --topic my-schema-registry-json Topic: my-schema-registry-json TopicId: 2xQXf5EXS-i96w-W-laIhw PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-schema-json Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: メッセージ送受信 送信側のターミナルで、Schema Registry を体験するための送信先のトピックを指定してプロデューサーを起動します。標準の Kafka CLI ツールは Schema Registry との連携に対応していないため、本章では Python で実装したオリジナルのプロデューサープログラムを使用します。 【送信側】 $ docker cp ~/handson/hands-on-kafka/try-my-hand/schema-registry-json-producer.py \ handson-client:/tmp/schema-registry-json-producer.py $ docker exec -i handson-client python3 /tmp/schema-registry-json-producer.py \ --bootstrap-server kafka:29092 \ --topic my-schema-registry-json \ --schema-registry-url http://schema-registry:8080/apis/ccompat/v7 \ --schema-id 1 > 受信側のターミナルで、Schema Registry を体験するための受信元のトピックを指定してコンシューマーを起動します。コンシューマー側でも、プロデューサーと同様に Python で実装したオリジナルのコンシューマープログラムを使用します。 【受信側】 $ docker cp ~/handson/hands-on-kafka/try-my-hand/schema-registry-json-consumer.py \ handson-client:/tmp/schema-registry-json-consumer.py $ docker exec -it handson-client python3 /tmp/schema-registry-json-consumer.py \ --bootstrap-server kafka:29092 \ --topic my-schema-registry-json \ --schema-registry-url http://schema-registry:8080/apis/ccompat/v7 \ --from-beginning \ --print-schema-ids 送信側のプロデューサーから JSON 形式のメッセージ送信します。 【送信側】 $ docker exec -i handson-client python3 /tmp/schema-registry-json-producer.py \ --bootstrap-server kafka:29092 \ --topic my-schema-registry-json \ --schema-registry-url http://schema-registry:8080/apis/ccompat/v7 \ --schema-id 1 > {"id":1, "name":"taro", "gender":"M", "age":10} スキーマと一致するため受信側のコンシューマーにメッセージが届きました。 【受信側】 $ docker exec -it handson-client python3 /tmp/schema-registry-json-consumer.py \ --bootstrap-server kafka:29092 \ --topic my-schema-registry-json \ --schema-registry-url http://schema-registry:8080/apis/ccompat/v7 \ --from-beginning \ --print-schema-ids schemaId:1 {"id": 1, "name": "taro", "gender": "M", "age": 10} スキーマと一致しないメッセージの場合は、エラーと判断されプロデューサーから送信できません。 【送信側】 $ docker exec -i handson-client python3 /tmp/schema-registry-json-producer.py \ --bootstrap-server kafka:29092 \ --topic my-schema-registry-json \ --schema-registry-url http://schema-registry:8080/apis/ccompat/v7 \ --schema-id 1 > {"id":1, "name":"taro", "gender":"M", "age":10} > {"id":2, "name":"unknown", "gender":"A", "age":20} [ERROR] Schema validation failed: 'A' is not one of ['M', 'F', 'X'] > {"id":3, "gender":"X", "age":30} [ERROR] Schema validation failed: 'name' is a required property > {"id":4, "name":"tamako", "gender":"F", "age":"forty"} [ERROR] Schema validation failed: 'forty' is not of type 'number' > {"id": 5, "name": "jiro", "gender": "M", "age": 50} id=1 は、スキーマと一致しているので、問題なくプロデューサーから送信された。 id=2 は、gender が規定値外であるため、エラーとしてプロデューサーから送信されない。 id=3 は、name フィールドが無いため、エラーとしてプロデューサーから送信されない。 id=4 は、age が数値型ではないため、エラーとしてプロデューサーから送信されない。 id=5 は、スキーマと一致しているので、問題なくプロデューサーから送信された。 受信側のコンシューマーでは、スキーマと一致するメッセージのみ受信できています。 【受信側】 $ docker exec -it handson-client python3 /tmp/schema-registry-json-consumer.py \ --bootstrap-server kafka:29092 \ --topic my-schema-registry-json \ --schema-registry-url http://schema-registry:8080/apis/ccompat/v7 \ --from-beginning \ --print-schema-ids schemaId:1 {"id": 1, "name": "taro", "gender": "M", "age": 10} schemaId:1 {"id": 5, "name": "jiro", "gender": "M", "age": 50} プロデューサーの起動とコンシューマーの起動は、簡単に実行できるように以下スクリプトを用意しています。 プロデューサーの起動: schema-registry-json-producer.sh コンシューマーの起動: schema-registry-json-consumer.sh Kafka Streams Application プロデューサーからコンシューマーへメッセージを送受信する過程で、プログラム言語で独自のデータ処理を実装できるストリームアプリケーションを利用したストリーム処理を体験します。 「my-streams-plaintext-input」トピックから受信して「my-streams-myhandson-output」トピックへ送信する過程で、Kafka Streams ライブラリを使用したアプリケーションでは以下の加工を行います。 「my-streams-plaintext-input」トピックから流れてきたテキスト形式のメッセージを1件ずつ、以下の加工を順次処理します。 StreamsBuilder.stream() メソッドを用いて「my-streams-plaintext-input」トピックに流れてきたメッセージを処理用のデータストリーム(KStream)として読み込みます。 流れてきたメッセージの値に対し、英大文字と小文字を相互に変換する加工処理を flatMapValues() で適用します。 加工後のデータストリームに対して to() メソッドを実行し、「my-streams-myhandson-output」トピックへ結果を書き込みます。 本ハンズオンでストリーム処理として実装した Kafka Streames Application プログラム: MyHandsOn.java まずは、「 Kafka 環境セットアップスクリプトの実行 」章で作成した、本章で利用するトピックが存在することを確認します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --describe \ --topic my-streams-plaintext-input,my-streams-myhandson-output Topic: my-streams-myhandson-output TopicId: dyClhMMiRxCWeZ4BebQzbg PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-streams-myhandson-output Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: Topic: my-streams-plaintext-input TopicId: RI97GYKxRu-LnAqZ1RZbQw PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-streams-plaintext-input Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: メッセージ送受信 送信側のターミナルで、Kafka Streams Application を体験するための送信先のトピックを指定してプロデューサーを起動します。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-streams-plaintext-input > 受信側のターミナルで、Kafka Streams Application を体験するための受信元のトピックを指定してコンシューマーを起動します。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-streams-myhandson-output \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.key=true \ --from-beginning 送信側のプロデューサーからメッセージ送信します。(Kafka Stream Application で加工前) 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-streams-plaintext-input > Hey, It's Me! 受信側のコンシューマーで、Kafka Stream Application で加工後のメッセージが届きました。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-streams-myhandson-output \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.key=true \ --from-beginning Partition:0 Offset:0 null hEY, iT'S mE! プロデューサーの起動とコンシューマーの起動は、簡単に実行できるように以下スクリプトを用意しています。 プロデューサーの起動: streams-producer.sh コンシューマーの起動: streams-consumer-myhandson.sh Kafka Streams と Flink DataStream API の使い分け ソースコードで実装する Kafka Streams と Flink DataStream API の違いは以下の通りです。 役割と構造の違い Kafka Streams(軽量な組み込みライブラリ) Spring Boot などの Java/Kotlin アプリケーション環境に Kafka Streams ライブラリを組み込むだけで、簡単に開発・動作環境を用意できます。Kafka からデータを受け取り、加工して、別の Kafka トピックへ戻す処理を外部クラスタなしでシンプルに構築できます。 Flink DataStream API(汎用ストリーム処理基盤) Flink クラスタ上で分散実行されるプログラムを書くための API です。Kafka 以外のデータソース(S3、データベース、他システム)とも自由につなぎ、巨大なデータを分散して高速かつ複雑に計算できます。 使い分けの判断基準 Kafka Streams を選ぶべきケース システム構成が Kafka 中心 で完結している。(Kafka ➔ 加工 ➔ Kafka) Flink のような 新しいクラスタを管理・運用したくない。(インフラ運用の負担を減らしたい) Java/Kotlin アプリの一機能として、手軽にストリーム処理を組み込みたい。 Flink DataStream API を選ぶべきケース Kafka だけでなく S3 や 外部 DB など多様なシステムを跨ぐデータパイプライン を作りたい。 複雑なデータ遅延を含む時間枠ごとの集計(ウィンドウ処理)を行いたい。(例:直近5分間の件数を1分ごとに集計) 複数の出来事が特定の順番・時間内で起きたパターンを検知したい。(CEP:例:ログイン失敗3回後に高額決済) 将来的にデータ量が膨大になり、クラスタレベルでの強力なスケールアウトを行いたい。 Kafka Connect Debezium Connector for MySQL (CDC) MySQL を対象とした CDC の Kafka Connector である「 Debezium connector for MySQL :: Debezium Documentation 」を利用して、データベースにコミットされたすべての操作をトレースしてみます。 AKHQ で登録されている Debezium Connector を確認できます。 http://localhost:8080/ui/hands-on-kafka/connect/debezium まずは、「 Kafka 環境セットアップスクリプトの実行 」章で作成した、本章で利用するトピックが存在することを確認します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 --describe \ --topic my-cdc-mysql.mytest.user Topic: my-cdc-mysql.mytest.user TopicId: W2iCPAiXR02_bcYvxmK_Nw PartitionCount: 1 ReplicationFactor: 1 Configs: min.insync.replicas=1 Topic: my-cdc-mysql.mytest.user Partition: 0 Leader: 1 Replicas: 1 Isr: 1 Elr: LastKnownElr: SQL 送信でメッセージ受信 先に受信側のターミナルで、CDC Connector を体験するための受信元のトピックを指定してコンシューマーを起動します。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-cdc-mysql.mytest.user \ --from-beginning 送信側のターミナルで、MySQL に対して更新系の SQL を送ります 【送信側】 $ docker exec mysql mysql -u myuser -pmypass -D mytest -e "INSERT INTO user(name) VALUES ('taro sios');" $ docker exec mysql mysql -u myuser -pmypass -D mytest -e "INSERT INTO user(name) VALUES ('hanako sios');" $ docker exec mysql mysql -u myuser -pmypass -D mytest -e "UPDATE user SET name='taro2 sios2' WHERE id=1;" $ docker exec mysql mysql -u myuser -pmypass -D mytest -e "DELETE FROM user WHERE id = 2;" コマンド実行で以下警告メッセージを表示する場合がありますが、ローカル環境の体験向けに、MySQL のログインパスワードを含んだコマンドラインを実行した警告であるため、今回は無視してください。 mysql: [Warning] Using a password on the command line interface can be insecure. 受信側のコンシューマーに、MySQL でデータが作成、更新、および削除された差分情報のメッセージが届きました。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-cdc-mysql.mytest.user \ --from-beginning {"before":null,"after":{"id":1,"name":"taro sios","updated_at":1786732821000}, ... } {"before":null,"after":{"id":2,"name":"hanako sios","updated_at":1786732835000}, ... } {"before":{"id":1,"name":"taro sios","updated_at":1786732821000},"after":{"id":1,"name":"taro2 sios2","updated_at":1786732849000}, ... } {"before":{"id":2,"name":"hanako sios","updated_at":1786732835000},"after":null, ... } 更新系 SQL の発行とコンシューマーの起動は、簡単に実行できるように以下スクリプトを用意しています。 更新系 SQL の発行: cdc-mysql-send-query.sh コンシューマーの起動: cdc-mysql-consumer.sh パーティション トピックをパーティションで分割することにより、コンシューマーの並列化による負荷分散を実現します。 パーティションを使わない場合 データの配信方法:すべてのメッセージが1つの単一レーン(トピック)に流れます。 コンシューマー側の挙動: 単一レーン構造:受信経路が1つのため、コンシューマーを追加しても並列分散処理を行う構成にはなりません。 独立した重複受信(マルチユース):追加した各コンシューマーが全データを受け取り、流れてきた注文情報に対して「発送処理」や「在庫管理」など別々の目的で利用できます。 パーティションを使った場合 データの配信方法:送信時にキーのハッシュ値を計算し、各パーティション(Partition 0〜n)へメッセージを自動分散します(例: usr_1 ➔ Partition 0、usr_2 ➔ Partition 1、usr_3 ➔ Partition 2)。 コンシューマー側の挙動: 理想的な並列処理(Group A):1つにつき1パーティションを専任(A-1〜A-3)して処理し、最高のパフォーマンスを発揮します。 余剰分の予備機化(Consumer A-4):パーティション数(3つ)を超える4個目は「ホットスタンバイ(待機状態)」となり、稼働中のコンシューマーが落ちた際に即座に引き継ぎます。 少ない数での兼任処理(Group B):コンシューマー数(2個)が少ない場合、1個(B-1)が複数パーティション(Partition 0 と 1)を自動で兼任して受信するため、データ欠落なく処理を実行します。 グループ間の独立配信:Group A と Group B のようにグループ ID を分けることで、全く同じデータ群をそれぞれの目的で独立して重複受信(マルチユース)できます。 メッセージ送受信 送信側のターミナルで、パーティションを体験するための送信先のトピックを指定してプロデューサーを起動します。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-topic-partition \ --reader-property "parse.key=true" --reader-property "key.separator=:" 受信側のターミナルで、パーティションを体験するための受信元のトピックを指定してコンシューマーを起動します。並列分散処理とホットスタンバイを確認するため、パーティション数(3つ)+ 予備(1個)= 計4つのターミナルを使い、それぞれで起動します。 【受信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic my-topic-partition \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.key=true \ --group my-group-a \ --from-beginning プロデューサーとコンシューマーのそれぞれが起動直後の状態であることを確認します。 送信側のプロデューサーから、 ”usr_1:” から ”usr_3:” までの3つのメッセージを送信します。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-topic-partition \ --reader-property "parse.key=true" --reader-property "key.separator=:" > usr_1:(p-0)taro > usr_2:(p-1)hanako > usr_3:(p-2)tama 受信側で各パーティションにバランシングされたコンシューマーが、それぞれのメッセージを受信できていることを確認します。 送信側のプロデューサーから、追加で ”usr_4:” から ”usr_6:” までの3つのメッセージを送信します。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-topic-partition \ --reader-property "parse.key=true" --reader-property "key.separator=:" > usr_1:(p-0)taro > usr_2:(p-1)hanako > usr_3:(p-2)tama > usr_4:(p-2)tamako > usr_5:(p-1)jiro > usr_6:(p-0)saki 受信側では引き続き、各パーティションと紐づいているコンシューマーにメッセージが届きます。 ここで「Partition: 1」の受信を担っているコンシューマーのターミナルで「Ctrl + C」を押し、コンシューマーを強制終了します。 送信側のプロデューサーから追加で、 ”usr_7:” から ”usr_9:” までの3つのメッセージを送ります。 【送信側】 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 --topic my-topic-partition \ --reader-property "parse.key=true" --reader-property "key.separator=:" > usr_1:(p-0)taro > usr_2:(p-1)hanako > usr_3:(p-2)tama > usr_4:(p-2)tamako > usr_5:(p-1)jiro > usr_6:(p-0)saki > usr_7:pochi(p-2) > usr_8:saburo(p-1) > usr_9:yoko(p-0) ホットスタンバイだったコンシューマーが「Partition 2」へ割り当てられ、メッセージを受信し始めます。このように、スタンバイの活性化と同時に、これまで「Partition 2」を担当していたコンシューマーが「Partition 1」の受信を引き継ぐなど、グループ全体でリバランス(再割り当て)が行われることもあります。 プロデューサーの起動とコンシューマーの起動は、簡単に実行できるように以下スクリプトを用意しています。 プロデューサーの起動: topic-partition-producer.sh コンシューマーの起動: topic-partition-consumer.sh 以上で、Apache Kafka を活用したデータパイプラインとストリーム処理のハンズオンはすべて終了です。お疲れ様でした! シンプルなメッセージの送受信から始まり、Flink や Kafka Streams によるリアルタイムなデータ加工、Schema Registry によるデータ品質管理、さらには Debezium による CDC 連携やパーティションによる負荷分散まで一通り体験していただきました。Kafka を中心としたエコシステムが連携し、リアルタイムデータ基盤としてどのように機能するのか、その「手応え」を掴んでいただけたなら幸いです。 環境のクリーンアップ Kafka のハンズオンが一通り終わったら、これまでに利用した環境をクリーンアップしてハンズオンを終了します。 コンテナ削除スクリプトの実行 ターミナルを用いて、以下のディレクトリへ移動します。 $ cd ~/handson/hands-on-kafka/container/ 環境構築時にも使用したスクリプトに down オプションを指定して実行して、Kafka 環境の各種コンテナを停止・削除します。 $ ./CREATE_CONTAINER.sh down スクリプトの実行が終わったら、 list オプションを付けてスクリプトを実行し、起動しているコンテナが存在しない(一覧に表示されない)ことを確認します。 $ ./CREATE_CONTAINER.sh list ### START: Show a list of container ########## CONTAINER ID IMAGE COMMAND CREATED STATUS PORTS NAMES 以上で環境のクリーンアップは完了です。お疲れ様でした。 続く Appendix では、このハンズオン環境を裏で支えている設定ファイルやスクリプトについて解説します。「どうやってこの環境を作ったのか詳しく知りたい!」という方は、ぜひこのままご覧ください。 Appendix docker-compose / Dockerfile の解説 flink/Dockerfile 「 Dockerfile 」は、Apache Flink コンテナ向けに以下の環境を構築するため、「docker-compose.yml」とは別に用意しています。 Kafka コネクタの事前配置:Flink から Kafka トピックへアクセスできるよう、Kafka Connector JAR を取得・配置します。 PyFlink 実行環境の構築:Python 環境のセットアップを行い、PyFlink ライブラリをインストールします。 streams/Dockerfile 「 Dockerfile 」は、Kafka Streams アプリケーションコンテナ向けに以下の環境を構築するため、「docker-compose.yml」とは別に用意しています。 ビルド・実行環境の構築:コンテナ起動時に Java アプリケーションをソースコードからビルド・実行できるよう、Maven および JDK 環境をセットアップします。 handson-client/Dockerfile 「 Dockerfile 」は、ハンズオンの各種操作を行うクライアントコンテナ向けに以下の環境を構築するため、「docker-compose.yml」とは別に用意しています。 Python 実行環境の構築:Schema Registry 章などで使用する Python 製のオリジナルプロデューサー・コンシューマーを実行できるよう、Python 環境と関連ライブラリをセットアップします。 コンテナ構築スクリプトの解説 「 コンテナ構築スクリプトの実行 」章で使用する「 CREATE_CONTAINER.sh 」スクリプトの処理内容について解説します。 コンテナの構築 用意した「docker-compose.yml」のコンテナ定義ファイルで一気にコンテナを構築します。 $ docker-compose \ -f docker-compose.yml \ up -d -V --remove-orphans Kafka 各コンテナの構成ファイル解説 AKHQ 設定ファイル application.yml AKHQ の Web UI で各コンテナの情報を表示できるよう、Kafka、Schema Registry、および Debezium(CDC Connector)の接続先を設定しています。 MySQL 設定ファイル .env-mysql MySQL のルートユーザーパスワードやデータベース名などの環境変数を指定します。 my.cnf 文字コードの設定や、CDC で必要となるバイナリーログの有効化を指定します。 init.sql ハンズオンで使用するデータベースやテーブルを構築するための初期化 DDL です。 Kafka 環境セットアップスクリプトの解説 「 Kafka 環境セットアップスクリプトの実行 」章で使用する「 SETUP_HANDS-ON.sh 」スクリプトの処理内容について解説します。 必須コマンドの存在確認 Kafka 環境のセットアップをスクリプトで行う際に必要となるコマンドがインストールされているかどうかを確認します。コマンドの存在が確認できなかった場合には、スクリプトは停止しますので、手作業でコマンドのインストールをお願いします。詳細な実行内容は以下実装を確認してください。 SETUP_HANDS-ON.sh トピックの作成 「handson-client」コンテナを介して、ハンズオンで利用するトピックを作成します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka:29092 --create --topic my-topic : 詳細な実行内容は以下実装を確認してください。 step01-create_topic.sh Flink の設定 Flink SQL のハンズオン向けに、事前に用意した DDL・DML を Flink SQL Client へ流し込みます。 $ cat ./config/flink_sql.sql | docker exec -i jobmanager ./bin/sql-client.sh flink_sql.sql Table API のハンズオン向けに、スクリプトの実行を開始します。 $ docker exec jobmanager ./bin/flink run -d --python /opt/flink/table_api.py DataStream API のハンズオン向けに、スクリプトの実行を開始します。 $ docker exec jobmanager ./bin/flink run -d --python /opt/flink/datastream_api.py 詳細な実行内容は以下実装を確認してください。 step02-register_flink.sh Schema Registry の設定 ハンズオン向けに、スキーマをSchema Registry へ登録します。 schema_registry.json 詳細な実行内容は以下実装を確認してください。 step03-register_schema_registry.sh Kafka Streams Application の設定 Streams Application と作ったトピックを紐づけるためにコンテナを再起動します。 $ docker compose -f ~/handson/hands-on-kafka/container/docker-compose.yml restart streams 詳細な実行内容は以下実装を確認してください。 step04-bind_stream_to_topic.sh Debezium Connector for MySQL の設定 Debezium Connector コンテナへ REST API で、MySQL をデータソースとした CDC を設定します。 $ curl -v -X POST \ "http://localhost:8083/connectors" \ -H "Content-Type: application/json" \ -d @$HOME/handson/hands-on-kafka/setup/config/debezium.json | \ jq 詳細な実行内容は以下実装を確認してください。 step05-register_connector.sh Kafka Streams Application の開発と起動 ハンズオン環境では、実装済みのソースコードからビルドしたパッケージを事前準備したコンテナで体験していただきましたが、ここでは Kafka Streams Application のソースコードをゼロから生成し、パッケージ作成から起動・メッセージ送受信までの手順を解説します。 まずは、「streams」コンテナを停止します。 $ docker stop streams Kafka Streams Application の開発と起動用の作業ディレクトリを用意します。 $ mkdir -p ~/work $ cd ~/work/ Kafka Streams Application の開発は Java 環境を必要とするため JDK をインストールします。 $ sudo apt install -y openjdk-21-jdk-headless Java21 以外のバージョンがアクティブな場合は「update-alternatives」を利用して Java21 をアクティブにします。(参照先: 初期環境構築: ユーティリティツール on Ubuntu | Java コマンド (JDK) ) ビルドに Maven を使うためインストールします。 $ sudo apt install -y maven Maven コマンドで Kafka Streams Application のソーステンプレートを作成します。 $ mvn archetype:generate \ -DinteractiveMode=false \ -DarchetypeGroupId=org.apache.kafka \ -DarchetypeArtifactId=streams-quickstart-java \ -DarchetypeVersion=3.9.0 \ -DgroupId=jp.sios \ -DartifactId=apisl.handson.kafka.streams \ -Dversion=0.1 \ -Dpackage=jp.sios.apisl.handson.kafka.streams $ cd apisl.handson.kafka.streams/ 「pom.xml」を修正。 @@ -31,6 +31,8 @@ <properties> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <kafka.version>3.9.0</kafka.version> <slf4j.version>1.7.36</slf4j.version> + <maven.compiler.source>21</maven.compiler.source> + <maven.compiler.target>21</maven.compiler.target> </properties> @@ -57,69 +59,16 @@ <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> - <version>3.1</version> + <version>3.13.0</version> <configuration> - <source>1.8</source> - <target>1.8</target> + <source>21</source> + <target>21</target> </configuration> </plugin> </plugins> <pluginManagement> <plugins> - <plugin> - <artifactId>maven-compiler-plugin</artifactId> - <configuration> - <source>1.8</source> - <target>1.8</target> - <compilerId>jdt</compilerId> - </configuration> - <dependencies> - <dependency> - <groupId>org.eclipse.tycho</groupId> - <artifactId>tycho-compiler-jdt</artifactId> - <version>0.21.0</version> - </dependency> - </dependencies> - </plugin> - <plugin> - <groupId>org.eclipse.m2e</groupId> - <artifactId>lifecycle-mapping</artifactId> - <version>1.0.0</version> - <configuration> - <lifecycleMappingMetadata> - <pluginExecutions> - <pluginExecution> - <pluginExecutionFilter> - <groupId>org.apache.maven.plugins</groupId> - <artifactId>maven-assembly-plugin</artifactId> - <versionRange>[2.4,)</versionRange> - <goals> - <goal>single</goal> - </goals> - </pluginExecutionFilter> - <action> - <ignore/> - </action> - </pluginExecution> - <pluginExecution> - <pluginExecutionFilter> - <groupId>org.apache.maven.plugins</groupId> - <artifactId>maven-compiler-plugin</artifactId> - <versionRange>[3.1,)</versionRange> - <goals> - <goal>testCompile</goal> - <goal>compile</goal> - </goals> - </pluginExecutionFilter> - <action> - <ignore/> - </action> - </pluginExecution> - </pluginExecutions> - </lifecycleMappingMetadata> - </configuration> - </plugin> </plugins> </pluginManagement> </build> LineSplit アプリケーションのソースコードを開き、入出力するトピック名をデフォルトから用意したトピック名に修正します。 $ vim src/main/java/jp/sios/apisl/handson/kafka/streams/LineSplit.java : - builder.<String, String>stream("streams-plaintext-input") + builder.<String, String>stream("my-streams-plaintext-input") .flatMapValues(value -> Arrays.asList(value.split("\\W+"))) - .to("streams-linesplit-output"); + .to("my-streams-linesplit-output"); : パッケージを作成し、LineSplit アプリケーションを起動します。 $ mvn clean package $ mvn exec:java -Dexec.mainClass=jp.sios.apisl.handson.kafka.streams.LineSplit 送信側のターミナルでメッセージを送信します。(LineSplit アプリは、「my-stream-plaintext-input」から「my-stream-linesplit-output」トピックへ単語に分割して転送します) $ docker exec -it handson-client /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server kafka:29092 \ --topic my-streams-plaintext-input >I-have-an-apple >I_have_an_apple >I have an apple 受信側のターミナルでメッセージを受信します。 $ docker exec -it handson-client /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 \ --topic my-streams-linesplit-output \ --formatter-property print.partition=true \ --formatter-property print.offset=true \ --formatter-property print.key=true \ --from-beginning Partition:0 Offset:0 null I Partition:0 Offset:1 null have Partition:0 Offset:2 null an Partition:0 Offset:3 null apple Partition:0 Offset:4 null I_have_an_apple Partition:0 Offset:5 null I Partition:0 Offset:6 null have Partition:0 Offset:7 null an Partition:0 Offset:8 null apple 最後にクリーンアップ作業として、Kafka Streams Application 用のターミナルに戻り「Ctrl + C」で起動を停止し、ソースコードを削除します。 $ cd ~/work/ $ rm -rf apisl.handson.kafka.streams/ 「streams」コンテナを再起動します。 $ docker start streams まとめ Apache Kafka を中心に、すべてオープンソース(OSS)のプロダクトを組み合わせ、基本的なメッセージ送受信からリアルタイムなデータ加工(Flink / Kafka Streams)、データ品質管理(Schema Registry)、外部 DB との CDC 連携(Debezium)、そしてパーティションによる負荷分散まで、リアルタイムデータパイプラインの一連の流れをハンズオン形式で体験していただきました。 今回のハンズオンで体験・学習できる主なポイントは以下の通りです。 オール OSS で揃うリアルタイムデータ基盤 Kafka Broker、AKHQ、Flink、Debezium などのオープンソースのみをコンテナで一括構築し、ローカル環境で手軽に実践的なイベントストリーミング環境を用意できること。 用途に応じた柔軟なストリーム処理 SQL感覚で扱える Flink SQL、コードで動的に定義する Table API、高度な状態管理が可能な DataStream API、そして軽量な組み込みライブラリである Kafka Streams など、ユースケースや開発スタイルに応じた多様な加工手法。 データ品質の確保とリアルタイム外部連携 Schema Registry によるスキーマ検証で不正データの混入を防ぐ品質管理と、Debezium コネクタを用いた MySQL からの CDC(変更データ捕獲)によるデータベース更新の即時ストリーム化。 パーティションによる分散処理と高可用性 パーティションとコンシューマーグループを活用した並列分散処理によるスケールアウトと、プロセス障害時におけるホットスタンバイ機への自動リバランス(引き継ぎ)の仕組み。 「イベントストリーミング」や「ストリーム処理」、「CDC」といった概念は、言葉やアーキテクチャ図だけで理解しようとすると難しく感じられがちですが、実際に手元でプロデューサーからメッセージを送信し、リアルタイムに加工・分散受信される様子を目の当たりにすることで、そのパワフルさや本質を実感していただけたのではないでしょうか。 本ハンズオンで使用したスクリプトや設定ファイルはすべて GitHub で公開していますので、環境構築の裏側の仕組みを解析してみたり、ご自身の開発アプリケーションと連携させてカスタマイズしてみたりと、リアルタイムデータ基盤実践の第一歩としてぜひご活用ください。 最後までお読みいただき、ありがとうございました。 ご覧いただきありがとうございます! この投稿はお役に立ちましたか? 役に立った 役に立たなかった 0人がこの投稿は役に立ったと言っています。 The post Apache Kafka で体験する『イベントストリーミング入門』 first appeared on SIOS Tech Lab .

動画

該当するコンテンツが見つかりませんでした

書籍