📣
TiDB Cloud Premium はパブリックプレビュー中です。エンタープライズワークロード向けの無制限のスケーリング、即時の弾力性、高度なセキュリティを提供します。このページは自動翻訳されたものです。原文はこちらからご覧ください。

TiCDC Debeziumプロトコル



TiCDC Debezium 、データベースの変更をキャプチャするためのツールです。キャプチャされたデータベースの変更はそれぞれ「イベント」と呼ばれるメッセージに変換され、Kafka に送信されます。v8.0.0以降、TiCDCはDebezium形式でTiDBの行データ変更(DMLイベント)をKafkaに直接送信することをサポートしているため、これまでDebeziumのMySQL統合を使用していたユーザーにとって、MySQLデータベースからの移行が簡素化されます。TiCDC v8.5.4-release.1新しい TiCDC アーキテクチャ)以降、TiCDC は Debezium 形式で DDL イベントと WATERMARK イベントを送信することもサポートしています。

Debeziumメッセージ形式を使用する

Kafkaをダウンストリームシンクとして使用する場合は、 sink-uri設定でprotocolフィールドをdebeziumに指定します。TiCDCはイベントに基づいてDebeziumメッセージをカプセル化し、TiDBデータ変更イベントをダウンストリームに送信します。

    新しい TiCDC アーキテクチャを使用するには、TiCDC 設定項目 newarchtrue に設定します。

    新しい TiCDC アーキテクチャでは、Debezium プロトコルは次の種類のイベントをサポートします。

    • DDL event: DDL 変更レコードを表します。アップストリームの DDL ステートメントが正常に実行された後、DDL イベントはすべての Message Queue (MQ) パーティションに送信されます。

    • DML event: 行データ変更レコードを表します。DML イベントは行変更が発生したときに送信されます。変更後の行に関する情報が含まれます。

    • WATERMARK event: 特定の時点を表します。この時点より前に受信したイベントが完全であることを示します。WATERMARK イベントは TiDB 拡張フィールドにのみ適用され、sink-urienable-tidb-extensiontrue に設定した場合に有効になります。

    従来の TiCDC アーキテクチャでは、Debezium プロトコルは行変更イベントのみをサポートし、DDL イベントと WATERMARK イベントは無視されます。行変更イベントは、行のデータ変更を表します。行が変更されると、行変更イベントが送信され、変更前後の行に関する情報が含まれます。WATERMARK イベントはテーブルのレプリケーションの進行状況を示し、ウォーターマークより前のすべてのイベントが下流に送信済みであることを示します。

    Debezium メッセージ形式を使用するための構成例は次のとおりです。

    cdc cli changefeed create --server=http://127.0.0.1:8300 --changefeed-id="kafka-debezium" --sink-uri="kafka://127.0.0.1:9092/topic-name?kafka-version=2.4.0&protocol=debezium"

    Debeziumの出力形式には、下流のコンシューマーが現在の行のデータ構造をより適切に理解できるように、現在の行のスキーマ情報が含まれています。スキーマ情報が不要なシナリオでは、changefeed設定ファイルでdebezium-disable-schemaパラメータをtrueまたはsink-uriに設定することで、スキーマ出力を無効にすることもできます。

    さらに、元の Debezium 形式には、TiDB の CommitTS の一意なトランザクション識別子などの重要なフィールドが含まれていません。データの整合性を確保するために、TiCDC は Debezium 形式に CommitTsClusterID の 2つのフィールドを追加し、TiDB データ変更の関連情報を識別します。

    メッセージ形式の定義

    このセクションでは、DDL イベント、DML イベント、および WATERMARK イベントのメッセージ形式について説明します。

    DDL イベント(新しい TiCDC アーキテクチャ)

    TiCDC は、キーと値の両方を Debezium 形式でエンコードして、DDL イベントを Kafka メッセージにエンコードします。

    キーフォーマット

    { "payload": { "databaseName": "test" }, "schema": { "type": "struct", "name": "io.debezium.connector.mysql.SchemaChangeKey", "optional": false, "version": 1, "fields": [ { "field": "databaseName", "optional": false, "type": "string" } ] } }

    キーのフィールドには、データベース名のみが含まれます。各フィールドの説明は以下のとおりです。

    フィールド名説明
    payloadJSONデータベース名に関する情報。
    schema.fieldsJSONpayload内の各フィールドの型情報。
    schema.typeStringフィールドのデータ型。
    schema.optionalBooleanフィールドがオプションかどうかを示します。trueの場合、フィールドはオプションです。
    schema.versionStringスキーマのバージョン。

    値の形式

    { "payload": { "source": { "version": "2.4.0.Final", "connector": "TiCDC", "name": "test_cluster", "ts_ms": 0, "snapshot": "false", "db": "test", "table": "table1", "server_id": 0, "gtid": null, "file": "", "pos": 0, "row": 0, "thread": 0, "query": null, "commit_ts": 1, "cluster_id": "test_cluster" }, "ts_ms": 1701326309000, "databaseName": "test", "schemaName": null, "ddl": "RENAME TABLE test.table1 to test.table2", "tableChanges": [ { "type": "ALTER", "id": "\"test\".\"table2\",\"test\".\"table1\"", "table": { "defaultCharsetName": "", "primaryKeyColumnNames": [ "id" ], "columns": [ { "name": "id", "jdbcType": 4, "nativeType": null, "comment": null, "defaultValueExpression": null, "enumValues": null, "typeName": "INT", "typeExpression": "INT", "charsetName": null, "length": 0, "scale": null, "position": 1, "optional": false, "autoIncremented": false, "generated": false } ], "comment": null } } ] }, "schema": { "optional": false, "type": "struct", "version": 1, "name": "io.debezium.connector.mysql.SchemaChangeValue", "fields": [ { "field": "source", "name": "io.debezium.connector.mysql.Source", "optional": false, "type": "struct", "fields": [ { "field": "version", "optional": false, "type": "string" }, { "field": "connector", "optional": false, "type": "string" }, { "field": "name", "optional": false, "type": "string" }, { "field": "ts_ms", "optional": false, "type": "int64" }, { "field": "snapshot", "optional": true, "type": "string", "parameters": { "allowed": "true,last,false,incremental" }, "default": "false", "name": "io.debezium.data.Enum", "version": 1 }, { "field": "db", "optional": false, "type": "string" }, { "field": "sequence", "optional": true, "type": "string" }, { "field": "table", "optional": true, "type": "string" }, { "field": "server_id", "optional": false, "type": "int64" }, { "field": "gtid", "optional": true, "type": "string" }, { "field": "file", "optional": false, "type": "string" }, { "field": "pos", "optional": false, "type": "int64" }, { "field": "row", "optional": false, "type": "int32" }, { "field": "thread", "optional": true, "type": "int64" }, { "field": "query", "optional": true, "type": "string" } ] }, { "field": "ts_ms", "optional": false, "type": "int64" }, { "field": "databaseName", "optional": true, "type": "string" }, { "field": "schemaName", "optional": true, "type": "string" }, { "field": "ddl", "optional": true, "type": "string" }, { "field": "tableChanges", "optional": false, "type": "array", "items": { "name": "io.debezium.connector.schema.Change", "optional": false, "type": "struct", "version": 1, "fields": [ { "field": "type", "optional": false, "type": "string" }, { "field": "id", "optional": false, "type": "string" }, { "field": "table", "optional": true, "type": "struct", "name": "io.debezium.connector.schema.Table", "version": 1, "fields": [ { "field": "defaultCharsetName", "optional": true, "type": "string" }, { "field": "primaryKeyColumnNames", "optional": true, "type": "array", "items": { "type": "string", "optional": false } }, { "field": "columns", "optional": false, "type": "array", "items": { "name": "io.debezium.connector.schema.Column", "optional": false, "type": "struct", "version": 1, "fields": [ { "field": "name", "optional": false, "type": "string" }, { "field": "jdbcType", "optional": false, "type": "int32" }, { "field": "nativeType", "optional": true, "type": "int32" }, { "field": "typeName", "optional": false, "type": "string" }, { "field": "typeExpression", "optional": true, "type": "string" }, { "field": "charsetName", "optional": true, "type": "string" }, { "field": "length", "optional": true, "type": "int32" }, { "field": "scale", "optional": true, "type": "int32" }, { "field": "position", "optional": false, "type": "int32" }, { "field": "optional", "optional": true, "type": "boolean" }, { "field": "autoIncremented", "optional": true, "type": "boolean" }, { "field": "generated", "optional": true, "type": "boolean" }, { "field": "comment", "optional": true, "type": "string" }, { "field": "defaultValueExpression", "optional": true, "type": "string" }, { "field": "enumValues", "optional": true, "type": "array", "items": { "type": "string", "optional": false } } ] } }, { "field": "comment", "optional": true, "type": "string" } ] } ] } } ] } }

    前述の JSON データの主要フィールドの説明は以下のとおりです。

    フィールド名説明
    payload.ts_msNumberTiCDC がこのメッセージを生成した時点のタイムスタンプ(ミリ秒)。
    payload.ddlStringDDL イベントの SQL ステートメント。
    payload.databaseNameStringイベントが発生したデータベースの名前。
    payload.source.commit_tsNumberイベントの CommitTs 値。
    payload.source.dbStringイベントが発生したデータベースの名前。
    payload.source.tableStringイベントが発生したテーブルの名前。
    payload.tableChangesArrayスキーマ変更後のテーブルスキーマ全体の構造化表現。tableChanges フィールドには、テーブルの各カラムのエントリを含む配列が含まれます。構造化表現は JSON または Avro 形式でデータを表すため、コンシューマーは DDL パーサーで事前処理しなくてもメッセージを簡単に読み取れます。
    payload.tableChanges.typeString変更の種類を示します。値は次のいずれかです。CREATE はテーブルが作成されたこと、ALTER はテーブルが変更されたこと、DROP はテーブルが削除されたことを示します。
    payload.tableChanges.idString作成、変更、または削除されたテーブルの完全識別子。テーブル名変更の場合、この識別子は <old><new> のテーブル名を連結したものです。
    payload.tableChanges.table.defaultCharsetNamestringイベントが発生したテーブルの文字セット。
    payload.tableChanges.table.primaryKeyColumnNamesstringテーブルの主キーを構成するカラムの一覧。
    payload.tableChanges.table.columnsArray変更されたテーブルの各カラムのメタデータ。
    payload.tableChanges.table.columns.nameStringカラム名。
    payload.tableChanges.table.columns.jdbcTypeNumberカラムの JDBC 型。
    payload.tableChanges.table.columns.commentStringカラムのコメント。
    payload.tableChanges.table.columns.defaultValueExpressionStringカラムのデフォルト値。
    payload.tableChanges.table.columns.enumValuesStringカラムの列挙値。形式は ['e1', 'e2'] です。
    payload.tableChanges.table.columns.charsetNameStringカラムの文字セット。
    payload.tableChanges.table.columns.lengthNumberカラムの長さ。
    payload.tableChanges.table.columns.scaleNumberカラムのスケール。
    payload.tableChanges.table.columns.positionNumberカラムの位置。
    payload.tableChanges.table.columns.optionalBooleanカラムがオプションかどうかを示します。true の場合、カラムはオプションです。
    schema.fieldsJSONpayload 内の各フィールドの型情報。変更されたテーブル内のカラムのスキーマ情報も含みます。
    schema.nameStringスキーマの名前。形式は "{cluster-name}.{schema-name}.{table-name}.SchemaChangeValue" です。
    schema.optionalBooleanフィールドがオプションかどうかを示します。true の場合、フィールドはオプションです。
    schema.typeStringフィールドのデータ型。

    DML イベント

    TiCDC は、キーと値の両方を Debezium 形式でエンコードして、DML イベントを Kafka メッセージにエンコードします。

    キーフォーマット

    { "payload": { "tiny": 1 }, "schema": { "fields": [ { "field":"tiny", "optional":true, "type":"int16" } ], "name": "test_cluster.test.table1.Key", "optional": false, "type":"struct" } }

    キーのフィールドには、主キーまたは一意インデックス列のみが含まれます。各フィールドの説明は以下のとおりです。

    フィールド名説明
    payloadJSON主キーまたは一意インデックス列に関する情報。各フィールドのキーと値は、それぞれカラム名とその現在の値を表します。
    schema.fieldsJSONpayload 内の各フィールドの型情報。変更前後の行データのスキーマ情報を含みます。
    schema.nameStringスキーマの名前。形式は "{cluster-name}.{schema-name}.{table-name}.Key" です。
    schema.optionalBooleanフィールドがオプションかどうかを示します。true の場合、フィールドはオプションです。
    schema.typeStringフィールドのデータ型。

    値の形式

    { "payload": { "source": { "version": "2.4.0.Final", "connector": "TiCDC", "name": "test_cluster", "ts_ms": 0, "snapshot": "false", "db": "test", "table": "table1", "server_id": 0, "gtid": null, "file": "", "pos": 0, "row": 0, "thread": 0, "query": null, "commit_ts": 1, "cluster_id": "test_cluster" }, "ts_ms": 1701326309000, "transaction": null, "op": "u", "before": { "tiny": 2 }, "after": { "tiny": 1 } }, "schema": { "type": "struct", "optional": false, "name": "test_cluster.test.table1.Envelope", "version": 1, "fields": [ { "type": "struct", "optional": true, "name": "test_cluster.test.table1.Value", "field": "before", "fields": [{ "type": "int16", "optional": true, "field": "tiny" }] }, { "type": "struct", "optional": true, "name": "test_cluster.test.table1.Value", "field": "after", "fields": [{ "type": "int16", "optional": true, "field": "tiny" }] }, { "type": "struct", "fields": [ { "type": "string", "optional": false, "field": "version" }, { "type": "string", "optional": false, "field": "connector" }, { "type": "string", "optional": false, "field": "name" }, { "type": "int64", "optional": false, "field": "ts_ms" }, { "type": "string", "optional": true, "name": "io.debezium.data.Enum", "version": 1, "parameters": { "allowed": "true,last,false,incremental" }, "default": "false", "field": "snapshot" }, { "type": "string", "optional": false, "field": "db" }, { "type": "string", "optional": true, "field": "sequence" }, { "type": "string", "optional": true, "field": "table" }, { "type": "int64", "optional": false, "field": "server_id" }, { "type": "string", "optional": true, "field": "gtid" }, { "type": "string", "optional": false, "field": "file" }, { "type": "int64", "optional": false, "field": "pos" }, { "type": "int32", "optional": false, "field": "row" }, { "type": "int64", "optional": true, "field": "thread" }, { "type": "string", "optional": true, "field": "query" } ], "optional": false, "name": "io.debezium.connector.mysql.Source", "field": "source" }, { "type": "string", "optional": false, "field": "op" }, { "type": "int64", "optional": true, "field": "ts_ms" }, { "type": "struct", "fields": [ { "type": "string", "optional": false, "field": "id" }, { "type": "int64", "optional": false, "field": "total_order" }, { "type": "int64", "optional": false, "field": "data_collection_order" } ], "optional": true, "name": "event.block", "version": 1, "field": "transaction" } ] } }

    前述のJSONデータの主要なフィールドの説明は以下のとおりです。

    フィールド名説明
    payload.opString変更イベントのタイプ。"c"INSERTイベント、 "u"UPDATEイベント、 "d"DELETEイベントを示します。
    payload.ts_msNumberTiCDC がこのメッセージを生成したときのタイムスタンプ (ミリ秒単位)。
    payload.beforeJSONステートメントの変更イベント前のデータ値。イベントが"c"の場合、フィールドbeforeの値はnullになります。
    payload.afterJSONステートメントの変更イベント後のデータ値。イベントが"d"の場合、フィールドafterの値はnullになります。
    payload.source.commit_tsNumberイベントのCommitTs値。
    payload.source.dbStringイベントが発生したデータベースの名前。
    payload.source.tableStringイベントが発生するテーブルの名前。
    schema.fieldsJSONペイロード内の各フィールドの型情報。変更前後の行データのスキーマ情報を含みます。
    schema.fields[1].fields[n].tidb_typeStringpayload.after 内の各カラムの TiDB 型。このフィールドは enable-tidb-extension = true の場合にのみ存在します。
    schema.nameStringスキーマの名前(形式は"{cluster-name}.{schema-name}.{table-name}.Envelope"
    schema.optionalBooleanフィールドがオプションかどうかを示します。 trueの場合、フィールドはオプションです。
    schema.typeStringフィールドのデータ型。

    WATERMARK イベント(新しい TiCDC アーキテクチャ)

    TiCDC は WATERMARK イベントを Kafka メッセージにエンコードし、キーと値の両方を Debezium 形式でエンコードします。

    キーフォーマット

    { "payload": {}, "schema": { "fields": [], "optional": false, "name": "test_cluster.watermark.Key", "type": "struct" } }

    フィールドの説明は以下のとおりです。

    フィールド名説明
    schema.nameStringスキーマの名前。形式は "{cluster-name}.watermark.Key" です。

    値の形式

    { "payload": { "source": { "version": "2.4.0.Final", "connector": "TiCDC", "name": "test_cluster", "ts_ms": 0, "snapshot": "false", "db": "", "table": "", "server_id": 0, "gtid": null, "file": "", "pos": 0, "row": 0, "thread": 0, "query": null, "commit_ts": 3, "cluster_id": "test_cluster" }, "op": "m", "ts_ms": 1701326309000, "transaction": null }, "schema": { "type": "struct", "optional": false, "name": "test_cluster.watermark.Envelope", "version": 1, "fields": [ { "type": "struct", "fields": [ { "type": "string", "optional": false, "field": "version" }, { "type": "string", "optional": false, "field": "connector" }, { "type": "string", "optional": false, "field": "name" }, { "type": "int64", "optional": false, "field": "ts_ms" }, { "type": "string", "optional": true, "name": "io.debezium.data.Enum", "version": 1, "parameters": { "allowed": "true,last,false,incremental" }, "default": "false", "field": "snapshot" }, { "type": "string", "optional": false, "field": "db" }, { "type": "string", "optional": true, "field": "sequence" }, { "type": "string", "optional": true, "field": "table" }, { "type": "int64", "optional": false, "field": "server_id" }, { "type": "string", "optional": true, "field": "gtid" }, { "type": "string", "optional": false, "field": "file" }, { "type": "int64", "optional": false, "field": "pos" }, { "type": "int32", "optional": false, "field": "row" }, { "type": "int64", "optional": true, "field": "thread" }, { "type": "string", "optional": true, "field": "query" } ], "optional": false, "name": "io.debezium.connector.mysql.Source", "field": "source" }, { "type": "string", "optional": false, "field": "op" }, { "type": "int64", "optional": true, "field": "ts_ms" }, { "type": "struct", "fields": [ { "type": "string", "optional": false, "field": "id" }, { "type": "int64", "optional": false, "field": "total_order" }, { "type": "int64", "optional": false, "field": "data_collection_order" } ], "optional": true, "name": "event.block", "version": 1, "field": "transaction" } ] } }

    前述の JSON データの主要なフィールドの説明は以下のとおりです。

    フィールド名説明
    payload.opString変更イベントのタイプ。"m" は watermark イベントを示します。
    payload.ts_msNumberTiCDC がこのメッセージを生成した時点のタイムスタンプ(ミリ秒単位)。
    payload.source.commit_tsNumberイベントの CommitTs 値。
    payload.source.dbStringイベントが発生するデータベースの名前。
    payload.source.tableStringイベントが発生するテーブルの名前。
    schema.fieldsJSONpayload 内の各フィールドの型情報。変更前後の行データのスキーマ情報を含みます。
    schema.nameStringスキーマの名前。形式は "{cluster-name}.watermark.Envelope" です。
    schema.optionalBooleanフィールドがオプションかどうかを示します。true の場合、そのフィールドはオプションです。
    schema.typeStringフィールドのデータ型。

    Data type mapping

    TiCDC Debeziumメッセージのデータ形式マッピングは基本的にDebeziumデータ型マッピングルールに準拠しており、これはMySQL用Debeziumコネクタのネイティブメッセージと概ね一致しています。ただし、一部のデータ型については、TiCDC DebeziumメッセージとDebeziumコネクタメッセージの間に以下の違いがあります。

    • 現在、TiDB は、GEOMETRY、LINESTRING、POLYGON、MULTIPOINT、MULTILINESTRING、MULTIPOLYGON、GEOMETRYCOLLECTION などの空間データ型をサポートしていません。

    • Varchar、String、VarString、TinyBlob、MediumBlob、BLOB、LongBlobなどの文字列型データ型の場合、列にBINARYフラグが付いている場合、TiCDCはBase64でエンコードした後、String型としてエンコードします。列にBINARYフラグが付いていない場合は、TiCDCは直接String型としてエンコードします。ネイティブDebeziumコネクタは、 binary.handling.modeに従って異なる方法でエンコードします。

    • TiCDCは、 DECIMALとNUMERIC含むDecimalデータ型をfloat64型で表現します。ネイティブのDebeziumコネクタは、データ型の精度に応じて、float32またはfloat64でエンコードします。

    • TiCDC は REAL を DOUBLE に変換し、長さが 1 の場合は BOOLEAN を TINYINT(1) に変換します。

    • TiCDC では、BLOB、TEXT、GEOMETRY、または JSON カラムにはデフォルト値がありません。

    • Debezium は FLOAT データ "5.61""5.610000133514404" に変換しますが、TiCDC は変換しません。

    • TiCDC は FLOAT の flen を誤って出力します tidb#57060

    • カラムの照合順序が "utf8_unicode_ci" で文字セットが null の場合、Debezium は charsetName"utf8mb4" に変換しますが、TiCDC は変換しません。

    • TiCDC は ENUM 要素内の \ をエスケープされた引用符として扱いますが、Debezium は扱いません。たとえば、TiCDC は ("c,\'d','g,''h") のような ENUM 要素を ('c,'d', 'g,''h') にエンコードします。

    • TiCDC は '1000-00-00 01:00:00.000' のような TIME のデフォルト値を "1000-00-00" に変換しますが、Debezium は変換しません。

    このページは役に立ちましたか?