1. クイックスタート: DebeziumとKafkaでMySQLのCDCパイプラインを5分で構築する
これまで、Elasticsearchや他のマイクロサービスへデータを同期する際、私たちのチームでは5分ごとにupdated_atをポーリングするcronジョブを設定するのが通例でした。この方法は、usersテーブルが1,200万行に達するまでは問題なく動作していました。しかしデータ量が増加すると、DBでI/Oボトルネックが頻発し、cronジョブが走るたびにCPU使用率が90%まで跳ね上がるようになったのです。CDC(Change Data Capture)は、まさにこの課題を根本から解決するために生まれました。テーブルへ直接クエリを発行する代わりに、MySQLのバイナリログ(binlog)からストリームを直接読み取ります。
以下は、MySQL 8.0、ZooKeeper、Apache Kafka、Debezium Kafka Connectを含む完全なdocker-compose.ymlファイルです。
version: '3.8'
services:
zookeeper:
image: quay.io/debezium/zookeeper:2.5
ports:
- "2181:2181"
kafka:
image: quay.io/debezium/kafka:2.5
ports:
- "9092:9092"
links:
- zookeeper
environment:
- ZOOKEEPER_CONNECT=zookeeper:2181
- KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092
mysql:
image: mysql:8.0
ports:
- "3306:3306"
environment:
- MYSQL_ROOT_PASSWORD=debezium
- MYSQL_DATABASE=inventory
- MYSQL_USER=mysqluser
- MYSQL_PASSWORD=mysqlpw
command: >
--server-id=223344
--log-bin=mysql-bin
--binlog-format=ROW
--binlog-row-image=FULL
--gtid-mode=ON
--enforce-gtid-consistency=ON
debezium:
image: quay.io/debezium/connect:2.5
ports:
- "8083:8083"
links:
- kafka
- mysql
environment:
- BOOTSTRAP_SERVERS=kafka:9092
- GROUP_ID=1
- CONFIG_STORAGE_TOPIC=my_connect_configs
- OFFSET_STORAGE_TOPIC=my_connect_offsets
- STATUS_STORAGE_TOPIC=my_connect_statuses
サービス群を起動します。
docker compose up -d
Kafka Connectの起動完了まで約30秒待機した後、ポート8083のREST API経由でMySQL Connectorを登録します。
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" \
http://localhost:8083/connectors/ -d '{
"name": "inventory-connector",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"tasks.max": "1",
"database.hostname": "mysql",
"database.port": "3306",
"database.user": "root",
"database.password": "debezium",
"database.server.id": "184054",
"topic.prefix": "dbserver1",
"database.include.list": "inventory",
"schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
"schema.history.internal.kafka.topic": "schema-changes.inventory"
}
}'
動作確認用のテーブルを作成し、テストレコードを挿入します。
CREATE TABLE inventory.users (
id INT AUTO_INCREMENT PRIMARY KEY,
name VARCHAR(255),
email VARCHAR(255),
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
INSERT INTO inventory.users (name, email) VALUES ('Nguyen Van A', '[email protected]');
ターミナルでKafkaトピックを購読し、CDCイベントがリアルタイム(レイテンシは通常15ms未満)で配信されることを確認します。
docker compose exec kafka /kafka/bin/kafka-console-consumer.sh \
--bootstrap-server kafka:9092 \
--topic dbserver1.inventory.users \
--from-beginning
2. 内部動作の仕組み
2.1. DebeziumはどのようにBinlogを読み取るのか?
Debeziumは本質的にMySQLレプリカとして動作します。MySQLサーバーに接続してbinlogストリームを継続的に要求・読み取り、バイナリペイロードを構造化されたJSONにパースします。すべてのデータ書き込み操作はイベントに変換され、対応するKafkaトピックへ直接送信されます。
2.2. 必須となる3つのBinlog設定フラグ
- binlog-format=ROW: Debeziumがステートメント単位ではなく行単位でデータを解析するために必須です。
STATEMENTに設定されている場合、ログには生のSQL構文しか記録されず、Debeziumは各カラムがどのように変更されたかを正確に把握できません。 - binlog-row-image=FULL: MySQLは変更前(before)と変更後(after)の全データフィールドを完全に記録します。これにより、後続のコンシューマーは元のDBに再クエリすることなく、完全なコンテキストを受け取ることができます。
- gtid-mode=ON: グローバルトランザクション識別子を有効にします。マスターとレプリカ間でフェイルオーバーが発生した場合でも、Debeziumはイベントの欠落や重複を起こさずに正確なbinlogオフセットを維持できます。
2.3. KafkaにおけるJSONメッセージの構造
Kafka上の各メッセージは、schemaとpayloadの2つの部分で構成されます。AvroやProtobufを使用してスキーマを圧縮しない限り、基本的にはpayload内のフィールドを直接扱うことになります。
before: コマンド実行前の行データ(INSERT操作の場合は常にnull)。after: 変更後の新しい行データ(DELETE操作の場合はnull)。op: 操作タイプコード(c: 作成/Insert、u: 更新/Update、d: 削除/Delete、r: 初期スナップショット読み取り/Read)。ts_ms: MySQLでイベントがコミットされた際のエポックミリ秒単位のタイムスタンプ。
3. 本番環境向けの詳細設定
3.1. Debezium専用ユーザーの作成
本番環境のコネクタにroot権限を付与することは絶対に避けてください。レプリケーションの読み取りに必要な最小限の権限を持つ専用ユーザーを作成する必要があります。
CREATE USER 'debezium_user'@'%' IDENTIFIED BY 'MatKhauPhucTap123!@#';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'debezium_user'@'%';
FLUSH PRIVILEGES;
3.2. 機密テーブルおよびカラムのフィルタリング
パスワードハッシュやクレジットカード番号などの機密データがKafkaトピックへ流出するのを防ぐため、コネクタレベルでホワイトリストおよびブラックリストを設定します。
{
"table.include.list": "inventory.users,inventory.orders",
"column.exclude.list": "inventory.users.password_hash,inventory.users.credit_card"
}
3.3. 大規模テーブル起動時におけるスナップショットの最適化
デフォルトのsnapshot.mode = initialでは、コネクタ起動時にテーブル全体に対してSELECTを実行します。2,000万行を超えるような大規模テーブルの場合、テーブルロックやメモリ枯渇を引き起こすリスクがあります。インクリメンタルスナップショットへの移行や、読み取りバッチサイズの調整を検討してください。
{
"snapshot.mode": "when_needed",
"snapshot.locking.mode": "none",
"snapshot.fetch.size": 10240
}
4. 6ヶ月間の本番運用で得られた実践ノウハウ
4.1. Binlog保持期間(Retention)の課題: コネクタ停止時のオフセット消失を防ぐ
私たちのチームが実際に経験した手痛い失敗談があります。週末の間にワーカーノードのディスク容量が枯渇し、Kafka Connectが2日間にわたって停止してしまいました。月曜日の朝に復旧を試みたものの、MySQL側で古いbinlogがすでにパージされていました。その結果、Debeziumはオフセットを見失い、最初からフルスナップショットを再実行せざるを得なくなり、後続のコンシューマー全体で深刻な遅延が発生しました。安全を確保するため、binlogの保持期間は最低でも3日〜7日に設定しておきましょう。
-- MySQL 8.0: binlogを7日間(604800秒)保持する設定
SET GLOBAL binlog_expire_logs_seconds = 604800;
4.2. DELETE処理時におけるTombstoneレコードの落とし穴
MySQLで行が削除されると、DebeziumはKafkaへ2つのメッセージを連続して送信します。1つ目は古いデータを含むop: "d"のメッセージ、そして直後に送信されるのがnullペイロードを持つ空のメッセージです(これはKafkaのログコンパクションを機能させるためのTombstoneレコードと呼ばれます)。バックエンド側のコンシューマーがnullペイロードを適切にハンドリングしていない場合、即座にNullPointerExceptionが発生してしまいます。
4.3. PrometheusとGrafanaによるリアルタイム遅延(Lag)モニタリング
コンシューマー側のデータ遅延が発生してから事態に気づくのでは遅すぎます。DebeziumからJMXメトリクスをPrometheusへエクスポートし、特にMilliSecondsBehindSourceメトリクスを注視してください。ピーク時間帯にこの数値が5000msを超えるようであれば、Kafka ConnectワーカーのCPU/RAM増強や、DBとKafka間のネットワーク帯域を見直すべき重要な兆候です。

