Ubuntu - TECH PLAY - TECH PLAY

TECH PLAY

Ubuntu

イベント

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

マガジン

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

技術ブログ

! 本記事に登場する IP アドレス・MAC アドレス・ホスト名・ポート番号などは、すべて架空のものに置き換えています。また、原因と検証結果を損なわない範囲で、システム構成や規模に関する数値を一部抽象化しています。 こんにちは。セーフィーの城山です。 本記事では、約2年にわたって断続的に発生していた「コンテナタスクの起動直後に限ってTCP接続が失敗する」という障害を、AIエージェントとともに再調査し、LinuxカーネルのARPキャッシュという意外な根本原因にたどり着いて解決するまでの一部始終を紹介します。 要約 コンテナ(Fargate)タスクの起動直後に限って、同一サブネットの
はじめに こんにちは、サイオステクノロジーの小沼 俊治です。 「理屈はいいから、まずは実際に 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 .
はじめに 前回は本シリーズの第二弾として、ペネトレーションテスト(Penetration Test)の基本概念と実施時の注意点を整理したうえで、HTB(Hack The Box)ラボ環境を用いてウェブディレクトリのブルートフォース攻撃を実施し、攻撃者の情報収集手法を確認した。さらに、その結果から防御側の観点で得られる知見をまとめた。まだ読んでいない方は、先にこちらを確認していただきたい。 今回は、ペネトレーションテスト実施の続きを解説する。 【要注意】 事前の了承を得ないまま勝手にペネトレーションテストを実施することは禁止である。 とても重要な内容であるため、上記の文言は今後の記

動画

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

書籍