
MySQL
イベント

マガジン
技術ブログ
はじめに こんにちは、サイオステクノロジーの小沼 俊治です。 「理屈はいいから、まずは実際に 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 .
本記事は 2026 年 4 月 15 日 に公開された「 Get to insights faster using Notebooks in Amazon SageMaker Unified Studio 」を翻訳したものです。 本記事では、Amazon SageMaker Unified Studio のノートブックがインフラの設定を簡素化し、インサイト獲得までの時間を短縮する仕組みを紹介します。住宅価格データの分析、スケーラブルなデータテーブルの作成、分散プロファイリング、機械学習 (ML) モデルの学習を、単一のノートブック環境で行う流れをご覧いただけます。 データサイエンティストやアナリストは、分析を始めるまでにインフラ設定と複数データソースの認証管理に何日も費やすことがしばしばあります。Amazon Simple Storage Service (Amazon S3)、Amazon Redshift、Snowflake、ローカルファイルにまたがるデータを扱う場合、認証設定の繰り返し、コンピュートスケーリングの手動判断、ツール切り替えの手間が積み重なり、インサイト獲得が遅れます。 Amazon SageMaker Unified Studio のノートブックは、12 以上のデータソースへの即時アクセス、ローカルから分散処理までのコンピュートスケーリング、AI によるコード生成を、単一のブラウザベース環境で提供します。ポリグロットプログラミング、マルチエンジンコンピュート、AI 支援による開発を活用し、疑問からインサイトまでの道のりを短縮する方法を解説します。 Amazon SageMaker Unified Studio のノートブックとは Amazon SageMaker Unified Studio のノートブックは、データ分析、探索、エンジニアリング、機械学習ワークフローのためのインタラクティブな環境を提供します。統合された 5 つの機能を備えています。 ポリグロットプログラミング : 同じノートブック環境で Python と SQL を相互に切り替えながらコードを書けます 統合されたデータアクセス : Amazon S3、AWS Glue Data Catalog、Apache Iceberg テーブル、Snowflake や BigQuery などのサードパーティソースに保存されたデータへ即座に接続できます ネイティブな可視化 : Python や SQL の結果から直接チャートを作成し、没入感のあるデータ分析ができます AI による開発支援 : SageMaker Data Agent を通じて自然言語プロンプトでコードを生成できます。データ分析、データサイエンス、ML タスク向けの的確なチャットインターフェイスを備えています 柔軟なコンピュート : ニーズの拡大に応じて、基本的なインスタンスから GPU 搭載環境までスケールできます アーキテクチャ 本セクションでは、ノートブックのアーキテクチャを説明します。ノートブックは、複数のコンピュートエンジン、多様なデータソース、AI アシスタントを統合したクラウドネイティブアーキテクチャにより、ブラウザベースのシンプルさでエンタープライズ規模の分析を実現します。 プレゼンテーション層 Amazon SageMaker Unified Studio 経由でノートブックインターフェイスにアクセスします。コード実行用のコードセル、ドキュメント用のマークダウンセル、チャートやテーブル用の可視化セルを備えた、使い慣れたインターフェイスで操作できます。 コンピュート層 専用のノートブックサーバーがカーネルのライフサイクルとセッション状態を管理します。主要コンポーネントには、コード補完用の Language Server、データサイエンスライブラリをプリロードした Python 3.11 ランタイム、Python、PySpark、SQL の実行を同じノートブック内で処理する Polyglot Kernel が含まれます。作成した各ノートブックは、永続的な Amazon Elastic Block Store (Amazon EBS) ストレージにバックアップされます。 実行層 ノートブックは複数の実行エンジンをサポートし、コードを最適な処理エンジンに自動的にルーティングします。インメモリ実行は、小規模なデータセットや迅速なプロトタイピングを処理します。Amazon Athena 経由の Apache Spark は、Spark Connect を介した大規模分析向けの分散処理を提供します。Amazon Athena (Trino)、Amazon Redshift、Snowflake、BigQuery へのネイティブ接続により、SQL クエリを処理します。 データ統合 AWS ネイティブ (Amazon S3、AWS Glue、Amazon Athena、Amazon Redshift) とサードパーティ (Snowflake、BigQuery、PostgreSQL、MySQL) を含む 12 以上のデータソースに統一的にアクセスできます。サポートされている最新のデータソースについては、 Connect to data sources を参照してください。 AI 層 SageMaker Data Agent は、2 つのモードで利用できます。複数ステップの分析ワークフロー向けの Agent Panel と、セル単位のコード生成に特化した Inline Assistance です。詳細は Accelerate context-aware data analysis and ML workflows with Amazon SageMaker Data Agent を参照してください。 作業を保護するため、セキュリティはアーキテクチャ全体に組み込まれています。データアクセスは AWS Identity and Access Management (AWS IAM) の権限に従います。ノートブックとエージェントは、利用が許可されたデータソースにのみアクセスできます。コンポーネント間の通信は暗号化チャネルを使用し、ノートブックのストレージは保存時に暗号化されます。AI エージェントには、破壊的な操作を防ぐガードレールが組み込まれており、コンプライアンスと監査のためにやり取りをログに記録します。 前提条件 始める前に、以下が必要です。 Amazon SageMaker Unified Studio のリソースを作成する権限を持つ AWS アカウント。必要な権限の詳細は、 Set up IAM-based domains を参照してください。 Python プログラミングと SQL クエリの基本的な知識 データ分析の概念と ML ワークフローの理解 サンプルの住宅データセットへのアクセス (ウォークスルー内で提供) ノートブックを使い始める まず、 Amazon SageMaker コンソール を開き、 Get started を選択します。 データとコンピュートへのアクセス権を持つ既存の AWS Identity and Access Management (AWS IAM) ロールを選択するか、新規ロールを作成するかを求められます。本ウォークスルーでは、 Create a new role を選択し、他のオプションはデフォルトのままにします。 Set up を選択します。環境の準備には数分かかります。 ユースケース 本記事では、ノートブックと SageMaker Data Agent を使って以下を実行します。 データセットの操作: サンプルデータセット housing.csv をアップロードし、データエクスプローラーで探索 ポリグロットプログラミング: DuckDB 経由で SQL によるデータフレームのクエリ AWS Glue によるマルチエンジンアクセス: AWS Glue テーブルを作成し、Athena SQL/Spark エンジンで分散処理を実行 高度な分析: Athena Spark でデータプロファイリング AI 支援による開発: Data Agent でプロファイリングと ML のコードを生成 ML ワークフロー: Random Forest モデルを学習し、結果を評価 まずインターフェイスを見て、コア機能を確認しましょう。 インターフェイスの理解 ノートブックのインターフェイスは、コード実行用のセルとドキュメント用のマークダウンという、使い慣れたノートブック規約に従っています。ノートブック内では、現在のプログラミング環境 (Python 3.11 など) とコンピュートプロファイルの仕様を確認できます。このインターフェイスでは以下ができます。 ファイルの参照、データカタログの探索、サードパーティ接続の管理を通じてデータにアクセス ノートブックのコンテキスト内で作成された変数を監視 ワークロード要件に応じて仮想 CPU と RAM を調整し、GPU インスタンスまでオンデマンドでコンピュートリソースをスケール Python パッケージを必要に応じてインストールおよび設定してパッケージを管理 データセットの操作 本ウォークスルーでは、 このページ からダウンロードできる housing.csv サンプルデータセットを使用します (リンク先ではファイル名が canvas-sample-housing.csv になっています)。左パネルの Files アイコンを選択し、 Local タブを選択します。 Local タブでノートブックに CSV ファイルをアップロードします。 ノートブックはデータ資産への即時アクセスを提供します。データエクスプローラーを使えば、AWS Glue Data Catalog、Amazon S3 テーブルカタログ、Amazon S3 バケット、設定済みのサードパーティ接続を参照できます。 3 点リーダーのオプションメニューを選択します。 Read as dataframe を選択し、ノートブックに挿入されたセルを実行して結果を確認します。 import pandas as pd <<df_csv_xxxx>> = pd.read_csv('housing.csv') <<df_csv_xxxx>> データフレームを返すと、ノートブックは自動的なデータプロファイリングを備えたリッチなテーブル形式で表示します。 ポリグロットプログラミング: Python と SQL を同時に ノートブックのもっとも強力な機能の 1 つが、Python と SQL の相互運用性です。データを Python のデータフレームに読み込んだ後、すぐに SQL でクエリできます。たとえば、海への近さ別に人口と世帯数の合計を計算するには、次のように実行します。 select sum(population) ,sum(households),ocean_proximity from<<df_csv_xxxx>> group byocean_proximity ノートブックのオートコンプリート機能はコンテキスト内のデータフレームを認識するため、SQL クエリを直感的に書けます。 この SQL クエリは DuckDB (インメモリの SQL データベースエンジン) 上で実行され、個別のインストールやサーバー保守は不要です。DuckDB は軽量な設計で、Python、Java、その他の環境に組み込めるため、迅速な対話型データ分析に適しています。分散処理が必要な場合は、データセット用の AWS Glue テーブルを作成した後、Apache Spark や Trino などのエンジンを使用できます。 データセット用の AWS Glue テーブルを作成する AWS Glue テーブルを作成すると、Amazon Athena SQL (Trino) や Amazon Athena Spark など、AWS Glue カタログ対応のさまざまなエンジンでデータセットをクエリできます。これらのエンジンは、それぞれのワークロード要件に対して最適な価格性能比を提供します。 まず AWS Glue データベースを作成します。ノートブックで新しいセルを作成するために SQL を選択し、 Amazon Athena (SQL) を選びます。 次の SQL を実行してデータベースを作成します: create database demo; 次に、データエクスプローラーに移動し、左上の +Add を選択して Create table を選びます。先ほど作成したデータベースを選択し、テーブル名を入力します。先ほど使用した housing.csv データセットファイルをアップロードします。サイドパネルで Next を選択して、テーブルを作成します。 次に、Amazon Athena SQL を使って新しいセルでサンプル SQL クエリを実行してみましょう。 select sum(population) , sum(households), ocean_proximity fromdemo.housing group by ocean_proximity Athena Spark による高度な機能 住宅価格を予測する ML モデルを構築する前に、データセットをさらに分析し、追加のインサイトを得るためにデータプロファイリングを実行しましょう。高度な探索には、ノートブック内で Amazon Athena Spark を使用できます。そのために、組み込みの Spark セッションを持つ新しい Python セルを作成します。次のコードを実行して Spark のバージョンを確認します。 # Verify Spark version spark.version SageMaker Data Agent によるデータプロファイリング 定型的なコードを手動で書く代わりに、組み込みの生成 AI 機能を活用できます。 プロンプト : 「Perform data profiling and create visualization for housing table」 AI アシスタントは、基本統計の算出、列単位のプロファイリング、データ型の分析、欠損値の検出を含む、包括的なプロファイリングコードを生成します。 エージェントは AWS Glue Data Catalog にアクセスし、housing テーブルの構造を把握したうえで、対象の列とデータ型に合わせたプロファイリングコードを生成しました。この文脈認識により、汎用的なコードスニペットを自環境向けに調整する際にありがちな試行錯誤が減ります。生成されたコードを確認して実行します。応答が速いため、分析を効率的に反復できます。 エラーが発生した場合は、次の図のように Fix with AI で解決できます。実行時にエラーが発生すると、 Fix with AI 機能がトレースバックを分析し、根本原因を診断して修正済みコードを生成するため、分析を止めずに進められます。 ML モデルの学習 次に、Data Agent を使って住宅価格を予測するモデルの学習コードを生成します。 プロンプト: 「Generate code to train a model that predicts housing prices. Use table housing.」 AI アシスタントは、次のようなエンドツーエンドのコードを生成します。 Amazon Athena Spark を使って AWS Glue カタログから住宅データを読み込み、pandas に変換 文字列型の列を数値化し、One-Hot エンコーディングでエンコードし、欠損値を除去 Random Forest モデルを学習して住宅価格の中央値を予測 モデル性能を評価 (RMSE、MAE、R-square) 予測に対する重要度が高い上位 10 個の特徴量を表示 この複数ステップのオーケストレーションにより、データアクセスからモデル評価までのワークフロー全体を任せられ、開発時間を数時間短縮できます。 エラーが発生した場合は、結果のトレースバックセクションにある Fix with AI で解決できます。 このワークフローでは、ノートブックの統合機能を紹介しました。ローカルからファイルをアップロードし、マルチエンジンアクセス用に AWS Glue テーブルを作成し、Amazon Athena Spark で分散プロファイリングを実行し、AI 支援による ML 開発で住宅価格を予測しました。これらすべてを、ツールを切り替えずに 1 つのノートブック環境で行いました。 主なメリットとベストプラクティス Amazon SageMaker Unified Studio のノートブックには、いくつかの利点があります。 インサイト獲得までの時間短縮 : 従来の環境では、分析を始める前に設定に数時間かかることもあります。ノートブックはこうした負荷を回避し、すぐに作業を始められます。 コラボレーションの向上 : 一貫した環境でノートブックを共有でき、再現性が確保され、「自分の環境では動く」問題が減ります。 複雑性の低減 : データソースや処理エンジンごとに別々のツールを使い分けるのではなく、複数のデータソースとコンピュートエンジンに 1 つのインターフェイスからアクセスできます。 AI による開発の加速 : タスク特化型のコードを生成し、的確な提案を受けられるため、繰り返しのコーディング作業にかかる時間を減らせます。 スケーラブルな性能 : メガバイトからペタバイト規模のデータセットを、適切なコンピュートリソースで処理できます。データ量の増加に応じてシステムが自動的にスケールします。 ベストプラクティス まずは小さめのインスタンスから始め、ニーズの拡大に合わせてスケールアップする形で、適切な コンピュートプロファイル で開始する。 繰り返し作業や複雑な操作には、自然言語プロンプトで AI アシスタントを活用する 。 大規模処理には Amazon Athena Spark、データウェアハウスには Amazon Redshift、その他特定のワークロードには専用エンジンを使うなど、 エンジンを戦略的に組み合わせる 。 マークダウンセルを使い、コードと並行して 作業内容を文書化する ことで、生きたドキュメントを作る。 複雑なワークフローを論理的なステップに分割し、 複数のセルで整理する ことで、可読性とデバッグ性を高める。 クリーンアップ 今後の課金を避けるため、本ウォークスルーで作成したリソースを削除します。 Amazon SageMaker Unified Studio コンソールで、Notebook ページに移動する ノートブックを削除する AWS Glue Data Catalog から demo データベースと housing テーブルを削除する 本ウォークスルーで作成した Amazon SageMaker Unified Studio ドメインを削除する 本ウォークスルー用に新しい IAM ロールを作成した場合は、IAM コンソールから削除する まとめ 本記事では、Amazon SageMaker Unified Studio のノートブックが作業効率を高め、インサイトをより速く提供する仕組みを紹介しました。使い慣れたノートブックインターフェイスに、エンタープライズ規模のコンピュート、マルチエンジンサポート、生成 AI 支援を組み合わせることで、チームはデータと AI のワークフローを効率化できます。 Python と SQL の統合、多様なデータソースへの即時アクセス、的確なコード生成機能により、ノートブックは現代のデータチームにとって価値あるツールになります。探索的データ分析、複雑なデータパイプラインの構築、ML モデルの学習を、単一の直感的な環境で必要な柔軟性とパワーを備えて実行できます。 使い始める準備はできましたか? Amazon SageMaker Unified Studio で最初のノートブックを作成 し、数分でデータ分析を始めましょう。 追加機能を探る: 季節分解と予測を用いた時系列分析ワークフロー テキスト分類や感情分析のための自然言語処理パイプライン ML モデルのバージョン管理のための Amazon SageMaker Model Registry との統合 ペタバイト規模の処理に対応する高度な Spark 最適化手法 詳細情報: Amazon SageMaker Unified Studio Administrator Guide Amazon SageMaker Unified Studio User Guide Notebooks in SageMaker Unified Studio Use the SageMaker Data Agent Pricing information 著者について Praveen Kumar Praveen Kumar は AWS の Principal Analytics Solutions Architect で、クラウドベースのサービスを用いた最新のデータおよび分析アプリケーションの設計、構築、実装に精通しています。関心領域はサーバーレステクノロジー、データガバナンス、データドリブンな AI アプリケーションです。 Majisha Namath Parambath Majisha Namath Parambath は Amazon SageMaker の Principal Engineer で、AWS で 10 年以上の経験を積んでいます。エージェント型システムを重視した包括的なデータ分析とインタラクティブな機械学習機能を提供する次世代サービス、Amazon SageMaker Unified Studio の重要な取り組みを主導しています。専門はシステム設計、アーキテクチャ、部門横断の実行で、エンタープライズ規模でのセキュリティ、パフォーマンス、信頼性に注力しています。エンジニアリング以外では、読書、料理、スキーを楽しんでいます。 Siddharth Gupta Siddharth Gupta は SageMaker の Unified Experiences で生成 AI を率いており、AI システムが複雑なタスクをユーザーに代わって自律的に実行する、エージェント型体験の推進に注力しています。以前は AWS でエッジ機械学習ソリューションを率いていました。彼の仕事は、開発者やデータサイエンティストが AI とやり取りする方法の改善、より直感的なデータ統合の実現、機械学習モデルを構築・デプロイするためのより良いツール作りに焦点を当てています。University of Illinois at Urbana-Champaign 出身で、Yahoo、Glassdoor、Twitch で豊富な経験を積んでいます。 LinkedIn で連絡できます。 この記事は Kiro が翻訳を担当し、Solutions Architect の Sotaro Hikita がレビューしました。
このブログ記事は、AWS ソリューションアーキテクトの多田慎也が執筆し、株式会社MonotaRO が監修しています。 はじめに 株式会社MonotaRO (モノタロウ)は、工具、部品、消耗品をはじめとする間接資材のインターネット通販を展開し、連結売上高 3,338 億円(2025 年 12 月期)、登録ユーザー数 1,100 万以上、取扱商品点数 2,800 万点以上を誇る事業者向け EC プラットフォームです。同社はこのたび、事業の根幹を支える基幹データベースをオンプレミスの MySQL から Amazon Aurora MySQL(以下、Aurora MySQL)へ移行し、 Amazon Aurora Global Database (以下、Aurora Global Database)によるマルチリージョン DR 構成とフルマネージド運用を実現しました。本ブログでは、移行の背景と課題、選定理由、移行アプローチ、そして得られた成果についてご紹介します。 背景と課題 MonotaRO の基幹データベースは、日々の受発注処理をはじめとする事業のコアシステムとして、大量のクエリを処理しています。同社はクラウド活用を推進してきましたが、基幹データベースがオンプレミスに残っていること、そしてその基盤となる MySQL 環境そのものに、複数の領域で課題を抱えていました。 オンプレミスの基幹データベースがクラウド移行のボトルネックに 基幹アプリケーションやバッチ処理は、基幹データベースに対して大量のクエリを発行しています。アプリケーションとデータベースの間にネットワーク遅延が入ると処理時間が大きく伸びるため、データベースがオンプレミスにある限り、関連するアプリケーション群もオンプレミスに残らざるを得ない状況でした。これにより、クラウドネイティブな開発基盤への移行が制約されていました。 さらに、基幹データベースは事業の広範な領域と密結合しており、ひとたび障害が発生すると、基幹業務のドメイン同士にとどまらず、EC サイトや周辺業務にまで影響が波及する構造でした。この「基幹データベースがオンプレにあるためにアプリもクラウドに移せない」「密結合ゆえに障害影響が広い」という課題は、多くの企業が直面する共通の課題です。 レプリケーション遅延の変動 オンプレミス環境では、プライマリとリードレプリカ間の論理レプリケーション遅延が処理内容によって変動し、状況によっては十数分規模に達することもありました。この遅延により、データ鮮度が求められる処理ではリードレプリカではなくプライマリを参照せざるを得ないケースが生じ、プライマリへの負荷増大を招く要因にもなっていました。 MySQL 保守期限への対応 基幹データベースは旧バージョンの MySQL で稼働しており、保守期限が迫っていました。セキュリティパッチの提供終了リスクへの対処が急務となっていました。 18 台レプリカのセルフマネージド運用負荷 参照負荷に対応するため 18 台のリードレプリカを運用していましたが、すべてセルフマネージドでした。スケールアップ/ダウン作業やメンテナンス対応など、多くの時間が割かれていました。 マネージドなマルチリージョン DR の不在 従来も障害時の復旧手段は備えていましたが、いずれも手動運用を前提としたセルフマネージドな構成であり、マネージドかつ自動化されたマルチリージョン DR は実現できていませんでした。事業の成長に伴い、リージョン障害にも耐えうる DR を、運用負荷を抑えたマネージドな形で整備し、事業継続性をさらに強化することが求められていました。 こうした課題に対し、Aurora の I/O 性能やストレージレベルレプリケーション機能の向上、大阪リージョンの活用によるレイテンシー面の改善、これまでの MySQL 移行で培った移行の練度といった追い風もあり、Aurora への移行が有力な選択肢として本格的に検討されました。 Aurora MySQL 選定の理由 MonotaRO が Aurora MySQL を選定した理由は、以下の 4 点に集約されます。 ストレージレベルレプリケーションによる遅延の根本解消 Aurora はストレージレベルでレプリケーションを行うため、従来の論理レプリケーションで発生していた変動の大きい遅延を根本的に解消できます。リーダー(Reader)は通常ミリ秒オーダーの遅延でライター(Writer)と同期されます。 MySQL 8.0 互換によるアプリケーション改修の最小化と EOSL 解消 Aurora MySQL は MySQL 8.0 互換であり、既存アプリケーションからの接続方法や SQL 構文の互換性が高いため、アプリケーション改修を最小限に抑えながらバージョンアップを実現できます。これにより、保守期限の問題もクラウド移行と同時に解決されます。 フルマネージドによる運用負荷削減と TCO 最適化 Aurora のマネージドサービスとしての特性により、18 台のセルフマネージドレプリカの運用負荷を大幅に削減できます。リージョン内でのインスタンス障害時の自動フェイルオーバー、自動バックアップ、パッチ適用など、従来手動で行っていた運用作業が自動化されます。さらに、ハードウェアの保守・更改やデータセンター運用が不要になり、必要な性能に応じてキャパシティを柔軟に調整できるため、運用工数まで含めた総所有コスト(TCO)の最適化にもつながります。 Aurora Global Database によるリージョンレベルの DR Aurora Global Database を採用することで、大阪リージョンをプライマリ、東京リージョンをセカンダリとするマルチリージョン構成を実現し、リージョンレベルの障害にも対応可能な DR 構成を構築できます。セカンダリリージョンへのレプリケーションは専用のストレージレベル基盤を通じて行われ、通常 1 秒未満の遅延で同期されます。リージョンレベルの切り替えは、スイッチオーバー(計画的な切り替え)や手動フェイルオーバーにより、迅速かつ制御された形で実施できます。 移行アプローチと技術的ポイント PoC による性能・信頼性の検証 MonotaRO の開発チームが主導し、約 4 ヶ月間の PoC(Proof of Concept)を実施しました。インフラストラクチャ観点での性能検証やフェイルオーバー動作の確認に加え、アプリケーションの動作検証も行い、本番ワークロードに耐えうることを確認した上で、Aurora MySQL の採用を正式に決定しています。 中継機を用いたレプリケーションでのデータ移行 MySQL のレプリケーションは隣接するメジャーバージョン間でのみサポートされるため、MySQL 5.6 から複数のメジャーバージョンをまたいで Aurora MySQL 3 系(MySQL 8.0 互換)へ直接レプリケーションすることはできません。そこで、MySQL 5.7 の中継用インスタンス(中継機)を挟み、「MySQL 5.6 →(中継機)MySQL 5.7 → Aurora MySQL 3 系」の順でバージョンを 1 段ずつ引き上げながらレプリケーションすることで、オンプレミスから Aurora へ安定してデータを同期しました(図 1)。 (図1: 中継機を用いた段階的レプリケーションによるデータ移行) 参照系から更新系への段階的な切り替え 移行はビッグバン的に一度で切り替えるのではなく、段階的なアプローチを採用しました。まず影響範囲の小さい参照系(Read)のアプリケーションから先に新しいデータベースのリーダーへ接続を切り替え、最後に更新系(Write)のアプリケーションを切り替えました。小さな単位で少しずつ切り替えることで、各段階で動作を検証しながら問題の早期発見と切り戻しを可能にし、ミッションクリティカルなシステムの移行リスクを最小化しました。 AWS Countdown Premium の活用 さらに MonotaRO は、こうした移行を確実に進めるため AWS Countdown Premium を活用しました。同社には過去にオンプレミスから Aurora へのデータベース移行経験がありましたが、今回は過去に対応したことのない規模であり、かつミッションクリティカルなシステムであることから、AWS の専門家による支援を組み合わせる判断をしました。 MonotaRO が AWS Countdown Premium の活用を通じて得た価値は、以下の通りです。 専門家レビューによるリスクの事前特定 : リソースレビューを通じて、現時点では顕在化していないものの、将来のスケールや運用の中で課題となり得る設計上のポイントを洗い出し、先回りして対処 インフラメトリクスの立ち合いレビュー : 移行前後のメトリクスを専門家がリアルタイムで確認し、開発チームと専門家の直接コミュニケーションにより、開発チームの安心感を醸成 包括的な視点でのアドバイス : データベースだけでなく、データベースに関連するネットワークや運用面も考慮した総合的なアドバイスを提供 以上のように、開発チーム主導の PoC による検証、参照系から更新系への段階的な切り替え、中継機を用いた安定したデータ移行、そして AWS Countdown Premium による専門家支援を組み合わせた結果、移行は計画メンテナンスの範囲内で実施され、予期しないサービス影響を出すことなく完了しました。 構成と成果 最終構成 移行後の構成は、Aurora Global Database を中核としたマルチリージョン構成です(図 2 参照)。 プライマリリージョン : 大阪リージョン セカンダリリージョン : 東京リージョン (図2: Aurora Global Database 構成図 ) 定量成果 指標 Before (オンプレミス MySQL) After (Aurora / Aurora Global Database) レプリケーション遅延 (DB 内リーダー同期) 処理内容により変動(状況により十数分規模) 通常 100 ミリ秒未満 ※1 運用対象レプリカ 18 台セルフマネージド フルマネージド DR 構成 手動運用のセルフマネージド構成 マネージドなマルチリージョン Global Database リージョン内フェイルオーバー (ライター障害)RTO 手動・約 2 時間 自動・1 分未満 ※2 ※1 本表のレプリケーション遅延は Aurora クラスター内のリーダー参照における値で、通常ミリ秒オーダー(多くの場合 100 ミリ秒未満)です。一方、Aurora Global Database によるクロスリージョンのレプリケーションは通常 1 秒未満です。なお、Aurora から下流のオンプレミス連携先へのデータ連携は、従来と同様の方式で行われます。 ※2 リージョン内フェイルオーバーは、障害の自動検知と自動フェイルオーバーによって実行されます。一般的な復旧時間は 60 秒未満(多くの場合 30 秒未満)であり、本記事では保守的に「1 分未満」と記載しています。 障害時の可用性向上と機会損失の削減 今回の移行で最も大きな効果は、 基幹データベース障害時の事業影響(機会損失)を大幅に削減できる 点です。 ライター(Writer)のダウン(AZ 障害)時の RTO は、従来の手動切り替えで約 2 時間を要していたところ、Aurora 化により自動フェイルオーバーで 1 分未満 に短縮されました。障害からの復旧時間が大幅に短縮されることで、ダウンタイムに伴う機会損失を削減できます。 リージョン障害時にも、Aurora Global Database のセカンダリリージョンへの昇格により、マネージドなクロスリージョン DR を実現しました。従来は手動運用を前提としたセルフマネージドな復旧構成であったのに対し、マネージド化・自動化されたことで運用負荷を抑えつつ、プライマリリージョンの復旧を待たずにサイトの注文受付を再開できます。 セキュリティ・運用面の向上 Aurora は AWS のマネージドサービスであり、ホスト OS への直接アクセスを前提としない運用となるため、実行環境に対する不正アクセスや改竄のリスクを低減できます。あわせて、パッチ適用やバックアップといった運用作業がマネージド化されることで、運用チームはより付加価値の高い業務に注力できるようになりました。 お客様の声 「当社はかねてより全社的なシステムのモダナイゼーションを進めてきましたが、その最大のボトルネックの1つであったオンプレの基幹データベースを Amazon Aurora 上へ移行させることができました。参照系から更新系へと段階的に切り替えることでリスクを下げ、業務影響を最小限にとどめながら安全にコントロールされた移行を実現できました。AWS様の的確なサポートにも大変助けられました。今回の移行は MySQL 8.0 へのバージョンアップを兼ねており、多数のリードレプリカに対するレプリケーション遅延の解消に加え、Instant DDL や Blue/Green Deployments によりスキーマ変更をかけにくいという長年の課題も解消。これによりデータ構造のリファクタリングが容易になりました。またクローン機能により検証用データベースも即時に用意できるようになり開発環境の自動化を加速できます。さらに CloudWatch Database Insights によるクエリパフォーマンスの可視化と改善サイクルを早くまわせる仕組みを導入準備中です。今後もモダナイゼーションを進めることでシステムの変更容易性を獲得し、当社の競争優位をさらに強固なものにしてまいります。」 — 株式会社MonotaRO 常務執行役 ソフトウェア本部長 普川氏 「オンプレミス環境の稼働メトリクスや全クエリを生成 AI も活用して事前に分析し、移行後に問題となり得る箇所を洗い出したうえで、参照系から更新系へと段階的に切り替えていきました。ミッションクリティカルなシステムでありながら、計画メンテナンスの範囲内で移行を完了できています。現在は、稼働状況をもとに Amazon Bedrock で改善案を生成し、プルリクエストとして開発チームへ提案する仕組みに加え、今後のバージョンアップに備えて現行バージョンと次期バージョンの挙動を突き合わせ、同様にプルリクエストとして提案する仕組みの構築を進めており、Aurora への移行はその土台となりました。」 — 株式会社MonotaRO ソフトウェア本部 コアシステムエンジニアリング部門 基幹モダナイゼーショングループ 石田氏(本プロジェクト PM) 「オンプレミスでは、データベース運用の多くが手作業であること、レプリケーション遅延が各システムに影響すること、データベースがオンプレミスにあるためアプリケーションもクラウドへ移せないことなど、長年にわたり様々な課題がありました。Aurora MySQL への移行でその多くに解決の見通しが立ち、システム設計と運用の自由度が大きく広がりました。」 — 株式会社MonotaRO ソフトウェア本部 プラットフォームエンジニアリング部門 サービスインフラグループ 中島氏(インフラ担当) 今後の展望 MonotaRO は以前からクラウド活用やモダナイゼーションを進めてきましたが、基幹データベースはオンプレミスに残り、関連するアプリケーション群のクラウド移行を妨げる要因となっていました。今回の移行によってこの制約が解消され、モダナイゼーションをさらに加速できる状態になりました。 基幹システムのモダナイゼーション加速 : これまで基幹データベースがオンプレミスに存在していたことで制約されていた関連アプリケーション群のクラウド移行が可能になりました。今後はクラウドネイティブ化(コンテナ基盤への移行)を進め、開発者が新サービス開発や事業成長に直結する業務により集中できる環境の実現を目指します。 リーダーへの読み取り振り分けによる負荷分散 : レプリケーション遅延が解消されたことで、これまでデータ鮮度の都合でライター(Writer)に向けざるを得なかった読み取り処理も、リーダー(Reader)へ振り分けられるようになります。ライターの負荷を軽減し、読み取りのスケールアウトを進める土台が整いました。 メンテナンス運用の最適化 : Aurora MySQL の Blue/Green Deployments を活用することで、メジャーバージョンアップやスキーマ変更といった計画メンテナンスを、ダウンタイムを最小化して実施できます。Blue/Green Deployments では、本番環境の複製(グリーン環境)で事前に変更を検証したうえで、通常 1 分未満・データ損失なしで切り替えが可能です。 継続的な運用改善 : Aurora が提供する CloudWatch Database Insights や Enhanced Monitoring などの機能を活用し、継続的な運用改善に取り組んでいます。 まとめ MonotaRO は、基幹データベースをオンプレミス MySQL から Aurora Global Database に移行し、レプリケーション遅延の根本解消、マルチリージョン DR の実現、18 台のセルフマネージドレプリカからフルマネージド運用への転換を達成しました。加えて、障害時の RTO 短縮により、障害時の事業影響(機会損失)を大幅に削減できる見込みです。この移行は、単なるデータベースの載せ替えではなく、クラウドネイティブな開発基盤への移行を加速させる重要な一歩です。














