このトピックでは、Kafka に書き込まれるメッセージの構造とフィールドについて説明します。
背景情報
データベース全体を Kafka に同期する場合、同期タスクは上流のデータソースからデータを読み取り、このトピックで説明する JSON フォーマットで Kafka トピックに書き込みます。メッセージの全体的なフォーマットには、変更記録のカラム情報と、変更前後のデータ状態が含まれます。コンシューマーがタスクの進捗を追跡できるように、同期タスクは定期的にハートビートレコードを Kafka トピックに送信します。このハートビートレコードには、値が MHEARTBEAT の op フィールドが含まれます。このトピックでは、メッセージの全体フォーマット、ハートビートメッセージフォーマット、およびソースデータ変更時のメッセージフォーマットについて説明します。詳細については、「フィールドの型」および「パラメーター」をご参照ください。
リアルタイム単一テーブル出力のメッセージフォーマット
リアルタイム単一テーブル同期タスクの Kafka 送信先を設定する際は、Kafka に書き込まれるレコードの値フォーマットを指定する必要があります。サポートされるフォーマットには、Canal CDC と JSON があります。詳細については、「付録:出力フォーマットの説明」をご参照ください。
フィールドの型
システムはソースからデータを読み取り、次の 6 種類のいずれか(BOOLEAN、DOUBLE、DATE、BYTES、LONG、STRING)にマッピングして、JSON フォーマットで Kafka トピックに書き込みます。
|
型 |
説明 |
|
BOOLEAN |
JSON のブール型に対応します。有効な値は |
|
DATE |
JSON の数値型に対応します。値は、ミリ秒 (ms) 精度の 13 桁の UNIX タイムスタンプです。 |
|
BYTES |
JSON の文字列型に対応します。Kafka に書き込まれる前に、バイト配列は Base64 エンコードされて文字列になります。消費時には、Base64 デコードする必要があります(エンコード:Base64.getEncoder().encodeToString(text.getBytes("UTF-8"))。デコード:Base64.getDecoder().decode(encodedText))。 |
|
STRING |
JSON の文字列型に対応します。 |
|
LONG |
JSON の数値型に対応します。 |
|
DOUBLE |
JSON の数値型に対応します。 |
パラメーター
以下の表は、Kafka に書き込まれるメッセージ内の各フィールドについて説明しています。
|
レベル 1 フィールド |
レベル 2 フィールド |
説明 |
|
schema |
dataColumn |
データ型は JSONArray です。dataColumn パラメーターは、上流のデータ変更レコードのすべてのカラムに関する型情報を記録します。変更操作には、データベース内のデータの追加・削除・更新などのデータ変更や、データベーステーブルスキーマの変更が含まれます。
|
|
primaryKey |
プライマリキー列名のリスト。 |
|
|
source |
ソースデータベースまたはテーブルに関する情報を含むオブジェクト。
|
|
|
payload |
before |
変更前のデータレコードを表す JSONObject。別名「before イメージ」とも呼ばれます。たとえば、MySQL データソースのレコードが更新された場合、
|
|
after |
変更後のデータ。別名「after イメージ」とも呼ばれます。形式は |
|
|
sequenceId |
StreamX が各レコードに対して生成する一意の文字列。完全データと増分データをマージする際にデータをソートするために使用されます。 説明
ソースからの更新操作により、2 つのレコード(「update before」レコードと「update after」レコード)が生成されます。両方のレコードは同じ |
|
|
scn |
ソースからのシステム変更番号(SCN)。このフィールドは Oracle ソースにのみ適用されます。 |
|
|
op |
ソースからキャプチャされた操作の種類。有効な値:
|
|
|
timestamp |
データレコードに関連するタイムスタンプを含む JSONObject。
|
|
|
ddl |
このフィールドは、データベースのテーブルスキーマが変更された場合にのみデータが入力されます。データの追加・削除・変更などのデータ変更では、対応する ddl は null に設定されます。
|
|
|
version |
N/A |
メッセージフォーマットのバージョン。 |
メッセージの全体フォーマット
以下のコードは、Kafka メッセージの全体フォーマットを示しています。
{
"schema": { // 変更に関するメタデータ(カラム名と型など)。
"dataColumn": [// 変更されたデータカラムの情報。送信先のレコード更新に使用。
{
"name": "id",
"type": "LONG"
},
{
"name": "name",
"type": "STRING"
},
{
"name": "binData",
"type": "BYTES"
},
{
"name": "ts",
"type": "DATE"
},
{
"name":"rowid",// データソースが Oracle の場合、rowid がデータカラムに含まれます。
"type":"STRING"
}
],
"primaryKey": [
"pkName1",
"pkName2"
],
"source": {
"dbType": "mysql",
"dbVersion": "1.0.0",
"dbName": "myDatabase",
"schemaName": "mySchema",
"tableName": "tableName"
}
},
"payload": {
"before": {
"dataColumn":{
"id": 111,
"name":"scooter",
"binData": "[base64 string]",
"ts": 1590315269000,
"rowid": "AAIUMPAAFAACxExAAE"// 文字列型。Oracle ソースからの rowid。
}
},
"after": {
"dataColumn":{
"id": 222,
"name":"donald",
"binData": "[base64 string]",
"ts": 1590315269000,
"rowid": "AAIUMPAAFAACxExAAE"// 文字列型。Oracle ソースからの rowid。
}
},
"sequenceId":"XXX",// 文字列型。完全データと増分データをマージする際にデータをソートするために使用。
"scn":"xxxx",// 文字列型。Oracle ソースからのシステム変更番号(SCN)。
"op": "INSERT/UPDATE_BEFOR/UPDATE_AFTER/DELETE/TRANSACTION_BEGIN/TRANSACTION_END/CREATE/ALTER/ERASE/QUERY/TRUNCATE/RENAME/CINDEX/DINDEX/GTID/XACOMMIT/XAROLLBACK/MHEARTBEAT...",// 値は大文字小文字を区別します。
"timestamp": {
"eventTime": 1,// 必須。ソースデータベースでの変更イベントの時刻。ミリ秒精度の 13 桁のタイムスタンプ。
"systemTime": 2,// 任意。同期タスクがこの変更メッセージを処理した時刻。ミリ秒精度の 13 桁のタイムスタンプ。
"checkpointTime": 3// 任意。同期オフセットをリセットする際に使用されるタイムスタンプ。通常は eventTime と同じ。
},
"ddl": {
"text": "ADD COLUMN ...",
"ddlMeta": "[シリアライズされた SQLStatement オブジェクトの Base64 エンコード文字列]"
}
},
"version":"1.0.0"
}
ハートビートメッセージフォーマット
{
"schema": {
"dataColumn": null,
"primaryKey": null,
"source": null
},
"payload": {
"before": null,
"after": null,
"sequenceId": null,
"timestamp": {
"eventTime": 1620457659000,
"checkpointTime": 1620457659000
},
"op": "MHEARTBEAT",
"ddl": null
},
"version": "0.0.1"
}
ソースデータ変更時のメッセージフォーマット
-
挿入操作のメッセージフォーマット:
{ "schema": { "dataColumn": [ { "name": "name", "type": "STRING" }, { "name": "job", "type": "STRING" }, { "name": "sex", "type": "STRING" }, { "name": "#alibaba_rds_row_id#", "type": "LONG" } ], "primaryKey": null, "source": { "dbType": "MySQL", "dbName": "pkset_test", "tableName": "pkset_test_no_pk" } }, "payload": { "before": null, "after": { "dataColumn": { "name": "name11", "job": "job11", "sex": "man", "#alibaba_rds_row_id#": 15 } }, "sequenceId": "1620457642589000000", "timestamp": { "eventTime": 1620457896000, "systemTime": 1620457896977, "checkpointTime": 1620457896000 }, "op": "INSERT", "ddl": null }, "version": "0.0.1" } -
更新操作のメッセージフォーマット:
-
「When one record in the source is updated, one Kafka record is generated.」オプションが選択されていない場合、更新操作により 2 つの Kafka メッセージが生成されます。1 つは更新前のデータ状態(「before イメージ」)、もう 1 つは更新後のデータ状態(「after イメージ」)です。以下のサンプルコードはそのフォーマットを示しています。
before イメージのメッセージフォーマット:
{ "schema": { "dataColumn": [ { "name": "name", "type": "STRING" }, { "name": "job", "type": "STRING" }, { "name": "sex", "type": "STRING" }, { "name": "#alibaba_rds_row_id#", "type": "LONG" } ], "primaryKey": null, "source": { "dbType": "MySQL", "dbName": "pkset_test", "tableName": "pkset_test_no_pk" } }, "payload": { "before": { "dataColumn": { "name": "name11", "job": "job11", "sex": "man", "#alibaba_rds_row_id#": 15 } }, "after": null, "sequenceId": "1620457642589000001", "timestamp": { "eventTime": 1620458077000, "systemTime": 1620458077779, "checkpointTime": 1620458077000 }, "op": "UPDATE_BEFOR", "ddl": null }, "version": "0.0.1" }after イメージのメッセージフォーマット:
{ "schema": { "dataColumn": [ { "name": "name", "type": "STRING" }, { "name": "job", "type": "STRING" }, { "name": "sex", "type": "STRING" }, { "name": "#alibaba_rds_row_id#", "type": "LONG" } ], "primaryKey": null, "source": { "dbType": "MySQL", "dbName": "pkset_test", "tableName": "pkset_test_no_pk" } }, "payload": { "before": null, "after": { "dataColumn": { "name": "name11", "job": "job11", "sex": "woman", "#alibaba_rds_row_id#": 15 } }, "sequenceId": "1620457642589000001", "timestamp": { "eventTime": 1620458077000, "systemTime": 1620458077779, "checkpointTime": 1620458077000 }, "op": "UPDATE_AFTER", "ddl": null }, "version": "0.0.1" } -
「When one record in the source is updated, one Kafka record is generated.」オプションが選択されている場合、更新操作により before イメージと after イメージの両方を含む 1 つの Kafka メッセージが生成されます。以下のサンプルコードはそのフォーマットを示しています。
{ "schema": { "dataColumn": [ { "name": "name", "type": "STRING" }, { "name": "job", "type": "STRING" }, { "name": "sex", "type": "STRING" }, { "name": "#alibaba_rds_row_id#", "type": "LONG" } ], "primaryKey": null, "source": { "dbType": "MySQL", "dbName": "pkset_test", "tableName": "pkset_test_no_pk" } }, "payload": { "before": { "dataColumn": { "name": "name11", "job": "job11", "sex": "man", "#alibaba_rds_row_id#": 15 } }, "after": { "dataColumn": { "name": "name11", "job": "job11", "sex": "woman", "#alibaba_rds_row_id#": 15 } }, "sequenceId": "1620457642589000001", "timestamp": { "eventTime": 1620458077000, "systemTime": 1620458077779, "checkpointTime": 1620458077000 }, "op": "UPDATE_AFTER", "ddl": null }, "version": "0.0.1" }
-
-
削除操作のメッセージフォーマット:
{ "schema": { "dataColumn": [ { "name": "name", "type": "STRING" }, { "name": "job", "type": "STRING" }, { "name": "sex", "type": "STRING" }, { "name": "#alibaba_rds_row_id#", "type": "LONG" } ], "primaryKey": null, "source": { "dbType": "MySQL", "dbName": "pkset_test", "tableName": "pkset_test_no_pk" } }, "payload": { "before": { "dataColumn": { "name": "name11", "job": "job11", "sex": "woman", "#alibaba_rds_row_id#": 15 } }, "after": null, "sequenceId": "1620457642589000002", "timestamp": { "eventTime": 1620458266000, "systemTime": 1620458266101, "checkpointTime": 1620458266000 }, "op": "DELETE", "ddl": null }, "version": "0.0.1" }