本記事では、シリアル化方式およびデータベースからテキストプロトコルへのデータ形式について説明します。
シリアライズ方式の形式説明
データ移行サービスを使用してソース側のデータをKafkaに移行する際、シリアライズ方式でターゲット側へのデータ移行メッセージ形式を制御できます。シリアライズ方式には、Default、Canal、DataWorks(V2.0以降をサポート)、SharePlex、DefaultExtendColumnType、Debezium、DebeziumFlatten、DebeziumSmt、およびAvroが含まれます。
説明
現在、シリアライズ方式 Debezium、DebeziumFlatten および DebeziumSmt は、OceanBaseデータベースのMySQL互換モードでのみサポートされています。
デフォルトのJSONメッセージ形式
データをKafkaに移行する際、デフォルトのシリアライズ方式で使用されるJSONメッセージ形式は以下のとおりです。
{
"prevStruct": { // 変更前のイメージ
"col1": "val1" // キー・バリュー・ペア。フルキー・バリューを含みます。
},
"postStruct": { // 変更後のイメージ
"col1": "val1" // キー・バリュー・ペア。フルキー・バリューを含みます。
},
"allMetaData"{
"checkpoint": "STRING", // 現在の同期ポイント。増分同期段階では同期された時間ポイント(秒単位のタイムスタンプ)を表し、フル移行段階では主キーのキー・バリュー・ペアを使用します。
"record_primary_key": "STRING", // 主キー列の名前。複数列が存在する場合は\u0001で区切ります。
"record_primary_value": "STRING", // 主キー値。複数列が存在する場合は\u0001で区切ります。
"source_identity": "STRING", // ソース側の識別子。増分同期の場合はサブトピック、フル移行の場合は意味のない番号です。
"dbType": "STRING", // データベースのタイプ。
"storeDataSequence": "LONG", // このフィールドは、増分同期シナリオでsource.jsonの設定にsequenceEnabled=trueが含まれている場合にのみ存在します。デフォルトはtrueです。ソートに使用され、生成規則は同期プロセス内で、タイムスタンプ + 5桁以下の連番が順次加算されることです。{タイムスタンプ}{連番}
"table_name": "STRING", // SQLステートメントを使用して変更するテーブル名
"db": "STRING", // SQLステートメントを使用して変更するデータベースの名前。OceanBaseデータベースの場合は、テナントを含み、形式は{tenant}.{database}です。
"timestamp": "STRING", // データ変更の秒単位のタイムスタンプ。増分同期にのみ存在します。
"uniqueId": "STRING", // 増分同期ではSTOREから受け渡されたトランザクション番号の識別子を表します。
"transId": "STRING", // OceanBaseデータベースではトランザクションIDを表します(incは自動インクリメント数、addrはコーディネーターアドレス、tsはトランザクションID生成時の時間、hashは上記3つのハッシュ値)。トランザクションが不完全な場合はnullです。
"clusterId": "STRING", // OceanBaseデータベースではclusterIdを表します。
"ddlType": "STRING", // DDLの具体的なタイプ
},
"recordType": "INSERT/UPDATE/DELETE/HEARTBEAT/DDL/ROW" // 変更タイプ
}
DDLのRecordでは、「ddl」のみが列名として存在し、値はDDLステートメントです。
前イメージと後イメージ:
prevStruct:増分データの前イメージ情報、すなわちSQL実行前のデータを表します。postStruct:増分データの後イメージ情報、すなわちSQL実行後のデータを表します。
DELETEにはprevStructのみが存在し、INSERTとDDLにはpostStructのみが存在します。UPDATEにはprevStructとpostStructの両方が存在し、HEARTBEAT(定期的なハートビートメッセージ)にはpostStructが存在しません。clusterIdの詳細な説明は以下のとおりです:OceanBaseデータベースでは、
ob_org_cluster_idはOceanBaseデータベースのセッションレベルのCLUSTER_IDを設定するために使用され、この値はトランザクションログに永続化されます。データ書き込み時にob_org_cluster_idを設定した場合は、設定した値が優先されます。それ以外の場合はすべてデフォルト値であり、増分データにのみ存在します。詳細については、ob_org_cluster_idをご参照ください。
データ例は以下のとおりです:
INSERT(挿入)データの例
{ "allMetaData":{ "checkpoint": null, "record_primary_key": "int8\u0001int16", "uniqueId": "{tid:11039xxxx127, partition_id:0, part_cnt:0},5917,391,0", "transId": "{hash:123456, inc:1234, addr:\"1.2.3.4:2883\", t:123456}", "clusterId": "123456", "source_identity": null, "record_primary_value": "3\u0001129", "dbType": "OB_MYSQL", "table_name": "table", "db": "tenant.database", "timestamp": "1609344671" }, "prevStruct": null, "recordType": "INSERT", "postStruct":{ "col1": 3, "col2": 129, "col3": 2147483646, "col4": 9223372036854775806, "col5": 10223372036854775806, "col6": 1.2222, "col7": 9.999999, "col8": "hello world", "col9": "aGVsbG8gd29ybGQ=", "col10": 9.99999999999, "col11": "2020-11-25", "col12": "00:01:02", "col13": "2020-11-25 00:01:02", "col14": "1606233662.012345", } }recordTypeがROWの場合、フルコンテンツで渡されるデータを表し、その形式はINSERTと同じです。{ "allMetaData":{ "checkpoint": null, "record_primary_key": "int8\u0001int16", "uniqueId": "{tid:11039xxxx127, partition_id:0, part_cnt:0},5917,391,0", "transId": "{hash:123456, inc:1234, addr:\"1.2.3.4:2883\", t:123456}", "clusterId": "123456", "source_identity": null, "record_primary_value": "3\u0001129", "dbType": "OB_MYSQL", "table_name": "table", "db": "tenant.database", "timestamp": "1609344671" }, "prevStruct": null, "recordType": "ROW", "postStruct":{ "col1": 3, "col2": 129, "col3": 2147483646, "col4": 9223372036854775806, "col5": 10223372036854775806, "col6": 1.2222, "col7": 9.999999, "col8": "hello world", "col9": "aGVsbG8gd29ybGQ=", "col10": 9.99999999999, "col11": "2020-11-25", "col12": "00:01:02", "col13": "2020-11-25 00:01:02", "col14": "1606233662.012345", } }データ更新(UPDATE)の例
{ "allMetaData": { "checkpoint": null, "record_primary_key": "int8\u0001int16", "source_identity": null, "uniqueId": "{tid:11039xxxx27, partition_id:0, part_cnt:0},5917,391,0", "transId": "{hash:123456, inc:1234, addr:\"1.2.3.4:2883\", t:123456}", "clusterId": "123456", "record_primary_value": "3\u0001129", "dbType": "OB_MYSQL", "table_name": "table", "db": "tenant.database", "timestamp": "1609344671" }, "prevStruct": { "col1": 3, "col2": 129, "col3": 2147483646, "col4": 9223372036854775806, "col5": 10223372036854775806, "col6": 1.2222, "col7": 9.999999999999, "col8": "hello world", "col9": "aGVsbG8gd29ybGQ=", "col10": 9.999999999999, "col11": "2020-11-25", "col12": "00:01:02", "col13": "2020-11-25 00:01:02", "col14": "1606233662.012345", }, "recordType": "UPDATE", "postStruct": { "col1": 3, "col2": 129, "col3": 2147483646, "col4": 9223372036854775806, "col5": 10223372036854775806, "col6": 1.2222, "col7": 9.999999999999, "col8": "hello world 2020", "col9": "aGVsbG8gd29ybGQ=", "col10": 9.999999999999, "col11": "2020-11-25", "col12": "00:01:02", "col13": "2020-11-25 00:01:02", "col14": "1606233662.012345", } }データ削除(DELETE)の例
{ "allMetaData":{ "checkpoint": null, "record_primary_key": "int8\u0001int16", "source_identity": null, "uniqueId": "{tid:11039xxxx27, partition_id:0, part_cnt:0},5917,391,0", "transId": "{hash:123456, inc:1234, addr:\"1.2.3.4:2883\", t:123456}", "clusterId": "123456", "record_primary_value": "3\u0001129", "dbType": "OB_MYSQL", "table_name": "table", "db": "tenant.database", "timestamp": "1609344671" }, "prevStruct":{ "col1": 3, "col2": 129, "col3": 2147483646, "col4": 9223372036854775806, "col5": 10223372036854775806, "col16": 1.2222, "col7": 9.99999999, "col8": "hello world", "col9": "aGVsbG8gd29ybGQ=", "col10": 9.999999999, "col11": "2020-11-25", "col12": "00:01:02", "col13": "2020-11-25 00:01:02", "col14": "1606233662.012345" }, "recordType": "DELETE", "postStruct": null }DDLの例
ALTER TABLE connector_test.all_mysql_type_test ADD column c90 VARCHAR(30) DEFAULT "test" COMMENT 'test';{ "prevStruct": null, "postStruct": { "ddl": "ALTER TABLE connector_test.all_mysql_type_test ADD column c90 VARCHAR(30) DEFAULT \"test\" COMMENT 'test'" }, "allMetaData": { "checkpoint": "1671177057", "dbType": "OB_MYSQL", "storeDataSequence": null, "db": "connector_test", "timestamp": "1671177057", "uniqueId": null, "ddlType": "ALTER_TABLE", "record_primary_key": null, "source_identity": null, "record_primary_value": null, "table_name": "all_mysql_type_test" }, "recordType": "DDL" }
Canal JSONメッセージフォーマット
データをKafkaに移行する際、Canalは以下のJSONメッセージ形式を使用してシリアライズします。
{
"database": "STRING", // SQLステートメントを使用して変更するデータベースの名前。OceanBaseデータベースの場合は、データベース名のみが存在し、テナント名は不要です。
"sqlType": {
"col1": "INTEGER", // 変更後の列タイプ。数値はjava.sql.Typesをご参照ください。
},
"data": [ // 変更後のデータキー・バリュー・ペア。現在は1つのメッセージのみが存在します。
{
"col1": "val1"
}
],
"pkNames": [ // 主キー列名
"col1"
],
"old": [ // 更新(UPDATE)メッセージにのみ存在します。UPDATEステートメントで変更される列、つまり変更前の列値を表します。
{
"col1": "val1"
}
],
"mysqlType": { // 列タイプの説明
"col": "STRING"
},
"type": "STRING", // 変更タイプ
"table": "STRING", // SQLステートメントを使用して変更するテーブル名
"es": "LONG", // 変更時間。ミリ秒単位のタイムスタンプ
"isDdl": "BOOLEAN", // DDLかどうか
"ts": "LONG", // 宛先への書き込みタイムスタンプ
"sql": "STRING" // 現在は空
}
データ例は以下のとおりです:
INSERT(挿入)データの例{ "database":"database", "sqlType":{ "col1":93, "col2":12, "col3":6, "col4":8, "col5":5, "col6":92, "col7":4, "col8":-5, "col9":2004, "col10":-6, "col11":91, "col12":3, "col13":-5, "col14":93 }, "data":[ { "col1":"2020-11-25 00:01:02", "col2":"hello world", "col3":1.2222, "col4":9.999999999999999093266253, "col5":129, "col6":"00:01:02", "col7":2147483646, "col8":9223372036854775806, "col9":"aGVsbG8gd29ybGQ=", "col10":3, "col11":"2020-11-25", "col12":9.9999999999999990, "col13":10223372036854775806, "col14":"1606233662.012345" } ], "pkNames":[ "col1", "col2" ], "old":null, "mysqlType":{ "col1":"datetime", "col2":"varchar", "col3":"float", "col4":"double", "col5":"smallint", "col6":"time", "col7":"int", "col8":"int64", "col9":"blob", "col10":"tinyint", "col11":"date", "col12":"decimal", "col13":"bigint", "col14":"timestamp" }, "type":"INSERT", "table":"table", "es":1609344671000, "isDdl":false, "ts":1618323429026, "sql":"" }UPDATE(更新)データの例{ "database":"database", "sqlType":{ "col1":93, "col2":12, "col3":6, "col4":8, "col5":5, "col6":92, "col7":4, "col8":-5, "col9":2004, "col10":-6, "col11":91, "col12":3, "col13":-5, "col14":93 }, "data":[ { "col1":"2020-11-25 00:01:02", "col2":"hello world 2020", "col3":1.2222, "col4":9.999999999999999093266253372484, "col5":129, "col6":"00:01:02", "col7":2147483646, "col8":9223372036854775806, "col9":"aGVsbG8gd29ybGQ=", "col10":3, "col11":"2020-11-25", "col12":9.9999999999999990932662, "col13":10223372036854775806, "col14":"1606233662.012345" } ], "pkNames":[ "col1", "col2" ], "old":[ { "string":"hello world" } ], "mysqlType":{ "col1":"datetime", "col2":"varchar", "col3":"float", "col4":"double", "col5":"smallint", "col6":"time", "col7":"int", "col8":"int64", "col9":"blob", "col10":"tinyint", "col11":"date", "col12":"decimal", "col13":"bigint", "col14":"timestamp" }, "type":"UPDATE", "table":"table", "es":1609344671000, "isDdl":false, "ts":1618364572908, "sql":"" }DELETE(削除)データの例{ "database":"database", "sqlType":{ "col1":93, "col2":12, "col3":6, "col4":8, "col5":5, "col6":92, "col7":4, "col8":-5, "col9":2004, "col10":-6, "col11":91, "col12":3, "col13":-5, "col14":93 }, "data":[ { "col1":"2020-11-25 00:01:02", "col2":"hello world", "col3":1.2222, "col4":9.99999999999999909326625, "col5":129, "col6":"00:01:02", "col7":2147483646, "col8":9223372036854775806, "col9":"aGVsbG8gd29ybGQ=", "col10":3, "col11":"2020-11-25", "col12":9.99999999999999909326625, "col13":10223372036854775806, "col14":"1606233662.012345" } ], "pkNames":[ "int8", "int16" ], "old":null, "mysqlType":{ "col1":"datetime", "col2":"varchar", "col3":"float", "col4":"double", "col5":"smallint", "col6":"time", "col7":"int", "col8":"int64", "col9":"blob", "col10":"tinyint", "col11":"date", "col12":"decimal", "col13":"bigint", "col14":"timestamp" }, "type":"DELETE", "table":"table", "es":1609344671000, "isDdl":false, "ts":1618364660278, "sql":"" }DDLの例
ALTER TABLE connector_test.all_mysql_type_test ADD COLUMN c90 VARCHAR(30) DEFAULT "test" COMMENT 'test';{ "database": "connector_test", "sqlType": null, "data": null, "pkNames": null, "old": null, "mysqlType": null, "type": "ALTER", "table": "all_mysql_type_test", "es": 1671177209000, "isDdl": true, "ts": 1671177291475, "sql": "ALTER TABLE connector_test.all_mysql_type_test ADD COLUMN c90 VARCHAR(30) DEFAULT \"test\" COMMENT 'test'" }
DataWorks JSONメッセージ形式
データをKafkaに移行する際、DataWorksが使用するシリアライズ方法は以下のJSONメッセージ形式です。
{
"version":"2.0", //プロトコルバージョン。現在はDataWorks 2.0バージョンのみサポートしています。
"schema": { //変更のメタデータ情報。列名と列型の情報のみ指定します。
"source": {//変更元の情報
"dbType": "mysql", //データソースタイプ
"dbVersion": "5.7.35", //データベースバージョン
"dbName": "myDatabase", //データベース名
"schema": "mySchema", //スキーマ名。スキーマが存在するシステムでは必須です。
"table": "tableName" //テーブル名
},
"column": [//変更されるデータ列の情報。更新対象テーブルのレコード内容
{
"name": "id",
"type": "bigint"
},
{
"name": "name",
"type": "varchar(20)"
},
{
"name": "mydata",
"type": "binary"
},
{
"name": "ts",
"type": "datetime"
}
],
"pk": [//主キーまたは一意キーがある場合は必須です。ない場合は省略可能です。
"pkName1",
"pkName2"
]
},
"payload": {
"before": {
"data":{
"id": 111,
"name":"scooter",
"mydata": "[base64 string]", //バイナリ型の場合は、Base64エンコードが必要です。
"ts": 1590315269000.123456789 //タイムスタンプ。整数部13桁、小数部9桁
}
},
"after": {
"data":{
"id": 222,
"name":"donald",
"mydata": "[base64 string]",
"ts": 1590315269000
}
},
"op":"INSERT/UPDATE/DELETE/HEARTBEAT/TRANSACTION_BEGIN/TRANSACTION_END/CREATE/ALTER/ERASE/QUERY/TRUNCATE/RENAME/CINDEX/DINDEX/GTID/XACOMMIT/XAROLLBACK/...",//大文字小文字を区別する
"timestamp": {
"eventTime": 1620457659000 // 変更がソースデータベースで発生した時間、ミリ秒単位の13桁のタイムスタンプ
},
"ddl": {
"text": "ADD COLUMN ..."
},
"scn": "⾃增 ID"
},
"extend": { //extend 拡張フィールド、今後の拡張要件に使用します。なければ、この部分は省略できます。
"load_fm":"CIBS" //出力元システムを記録します。例:"CIBS"
}
}
移行タスクのハートビートメッセージ:
{
"version": "2.0", //プロトコルバージョン
"payload": {
"timestamp": {
"eventTime": 1620457659000 //ハートビートパケットの時間
},
"op": "HEARTBEAT" //ハートビートパケットであることを示す
}
}
データの例は以下の通りです:
INSERT(挿入)データの例{ "version":"2.0", "schema":{ "source":{ "dbType":"ob_mysql", "dbVersion":null, "dbName":"db", "schema":null, "table":"tab" }, "column":[ { "name":"int8", "type":"TINYINT" }, { "name":"int16", "type":"SMALLINT" }, { "name":"int32", "type":"INT" }, { "name":"int64", "type":"INT64" }, { "name":"float32", "type":"FLOAT" }, { "name":"float64", "type":"DOUBLE" }, { "name":"bigInt", "type":"BIGINT" }, { "name":"boolean", "type":"BOOLEAN" }, { "name":"string", "type":"VARCHAR" }, { "name":"bytes", "type":"BLOB" }, { "name":"decimal", "type":"DECIMAL" }, { "name":"localDate", "type":"DATE" }, { "name":"localTime", "type":"TIME" }, { "name":"localDateTime", "type":"DATETIME" }, { "name":"timestamp", "type":"TIMESTAMP" }, { "name":"zonedDateTime", "type":"ZONED_DATETIME" }, { "name":"intervalDayToSecond", "type":"INTERVAL_DAY_TO_SECOND" }, { "name":"intervalYearToMonth", "type":"INTERVAL_YEAR_TO_MONTH" } ], "pk":[ "pkName1", "pkName12" ] }, "payload":{ "before":null, "after":{ "data":{ "col1":3, "col2":129, "col3":2147483646, "col4":9223372036854775806, "col5":1.2222, "col6":0.0000000000000008125, "col7":10223372036854775806, "col8":1, "col9":"hello world", "col10":"aGVsbG8gd29ybGQ=", "col11":0.000000000000125, "col12":"2020-11-25", "col13":"00:01:02", "col14":"2020-11-25 00:01:02", "col15":"1606233662.012345", "col16":"2020-11-25 00:01:02.012345 Asia/Shanghai", "col17":"INTERVAL '3' DAY", "col18":"INTERVAL '4' YEAR" } }, "op":"INSERT", "timestamp":{ "eventTime":1647581000000, "systemTime":1647581000795, "checkpointTime":1647581000 }, "ddl":null, "scn":"null" }, "extend":{ "load_fm": "test" } }UPDATE(更新)データの例{ "version":"2.0", "schema":{ "source":{ "dbType":"ob_mysql", "dbVersion":null, "dbName":"db", "schema":null, "table":"tab" }, "column":[ { "name":"int8", "type":"TINYINT" }, { "name":"int16", "type":"SMALLINT" }, { "name":"int32", "type":"INT" }, { "name":"int64", "type":"INT64" }, { "name":"float32", "type":"FLOAT" }, { "name":"float64", "type":"DOUBLE" }, { "name":"bigInt", "type":"BIGINT" }, { "name":"boolean", "type":"BOOLEAN" }, { "name":"string", "type":"VARCHAR" }, { "name":"bytes", "type":"BLOB" }, { "name":"decimal", "type":"DECIMAL" }, { "name":"localDate", "type":"DATE" }, { "name":"localTime", "type":"TIME" }, { "name":"localDateTime", "type":"DATETIME" }, { "name":"timestamp", "type":"TIMESTAMP" }, { "name":"zonedDateTime", "type":"ZONED_DATETIME" }, { "name":"intervalDayToSecond", "type":"INTERVAL_DAY_TO_SECOND" }, { "name":"intervalYearToMonth", "type":"INTERVAL_YEAR_TO_MONTH" } ], "pk":[ "pkName1", "pkName2" ] }, "payload":{ "before":{ "data":{ "col1":3, "col2":129, "col3":2147483646, "col4":9223372036854775806, "col5":1.2222, "col6":0.000000000000000125, "col7":10223372036854775806, "col8":1, "col9":"hello world", "col10":"aGVsbG8gd29ybGQ=", "col11":0.0000000000125, "col12":"2020-11-25", "col13":"00:01:02", "col14":"2020-11-25 00:01:02", "col15":"1606233662.012345", "col16":"2020-11-25 00:01:02.012345 Asia/Shanghai", "col17":"INTERVAL '3' DAY", "col18":"INTERVAL '4' YEAR" } }, "after":{ "data":{ "col1":3, "col2":129, "col3":2147483646, "col4":9223372036854775806, "col5":1.2222, "col6":0.00000000008125, "col7":10223372036854775806, "col8":1, "col9":"hello world 2020", "col10":"aGVsbG8gd29ybGQ=", "col11":0.000000000125, "col12":"2020-11-25", "col13":"00:01:02", "col14":"2020-11-25 00:01:02", "col15":"1606233662.012345", "col16":"2020-11-25 00:01:02.012345 Asia/Shanghai", "col17":"INTERVAL '3' DAY", "col18":"INTERVAL '4' YEAR" } }, "op":"UPDATE", "timestamp":{ "eventTime":1647581038000, "systemTime":1647581038674, "checkpointTime":1647581038 }, "ddl":null, "scn":"null" }, "extend":{ "load_fm": "test" } }DELETE(削除)データの例{ "version":"2.0", "schema":{ "source":{ "dbType":"ob_mysql", "dbVersion":null, "dbName":"db", "schema":null, "table":"tab" }, "column":[ { "name":"int8", "type":"TINYINT" }, { "name":"int16", "type":"SMALLINT" }, { "name":"int32", "type":"INT" }, { "name":"int64", "type":"INT64" }, { "name":"float32", "type":"FLOAT" }, { "name":"float64", "type":"DOUBLE" }, { "name":"bigInt", "type":"BIGINT" }, { "name":"boolean", "type":"BOOLEAN" }, { "name":"string", "type":"VARCHAR" }, { "name":"bytes", "type":"BLOB" }, { "name":"decimal", "type":"DECIMAL" }, { "name":"localDate", "type":"DATE" }, { "name":"localTime", "type":"TIME" }, { "name":"localDateTime", "type":"DATETIME" }, { "name":"timestamp", "type":"TIMESTAMP" }, { "name":"zonedDateTime", "type":"ZONED_DATETIME" }, { "name":"intervalDayToSecond", "type":"INTERVAL_DAY_TO_SECOND" }, { "name":"intervalYearToMonth", "type":"INTERVAL_YEAR_TO_MONTH" } ], "pk":[ "pkName1", "pkName2" ] }, "payload":{ "before":{ "data":{ "col1":3, "col2":129, "col3":2147483646, "col4":9223372036854775806, "col5":1.2222, "col6":0.0000000000125, "col7":10223372036854775806, "col8":1, "col9":"hello world", "col10":"aGVsbG8gd29ybGQ=", "col11":0.0000000000125, "col12":"2020-11-25", "col13":"00:01:02", "col14":"2020-11-25 00:01:02", "col15":"1606233662.012345", "col16":"2020-11-25 00:01:02.012345 Asia/Shanghai", "col17":"INTERVAL '3' DAY", "col18":"INTERVAL '4' YEAR" } }, "after":null, "op":"DELETE", "timestamp":{ "eventTime":1647581072000, "systemTime":1647581072976, "checkpointTime":1647581072 }, "ddl":null, "scn":"null" }, "extend":{ "load_fm": "test" } }DDLの例
ALTER TABLE connector_test.all_mysql_type_test ADD COLUMN c90 VARCHAR(30) DEFAULT "test" COMMENT 'test';{ "version": "2.0", "schema": { "source": { "dbType": "ob_mysql", "dbVersion": null, "dbName": "connector_test", "schema": null, "table": "all_mysql_type_test" }, "column": null, "pk": null }, "payload": { "before": null, "after": null, "op": "ALTER", "timestamp": { "eventTime": 1671177209000, "systemTime": 1671177291485, "checkpointTime": 1671177200 }, "ddl": { "text": "ALTER TABLE connector_test.all_mysql_type_test ADD COLUMN c90 VARCHAR(30) DEFAULT \"test\" COMMENT 'test'" }, "scn": "null" }, "extend": {} }
SharePlex JSONメッセージ形式
データをKafkaに移行する際、SharePlexが使用するシリアル化方式は以下のJSONメッセージ形式です。
{
"data": { // 変更データのキーと値のペア。INSERT/DELETEの場合は全量値、UPDATEの場合は変更値のみ
"col1": "val1"
},
"meta": {
"time": "YYYY-MM-DDTHH:mm:ss", // 変更時間
"op": "", // 変更タイプ。ins/upd/del/ddlを含む
"posttime": "YYYY-MM-DDTHH:mm:ss", // ターゲット側への書き込み時間
"idx": "STRING", //トランザクション内のメッセージのインデックス/インデックスされたメッセージ数。このパラメータは廃止されました。
"size": "NUMBER", //トランザクション内のメッセージ数。このパラメータは廃止されました。
"seq": "STRING", //ソート番号。ソース側でtransactionEnabledを有効にしている必要があります。
"table": "STRING", //SQL変更データベーステーブル名 {database}.{table}
"rowid": "STRING", // {変更データベーステーブル名}-{主キー値は\u0001}で分割
"trans": "STRING", //トランザクションID
"scn": "STRING" //このフィールドは、増分シナリオでsource.json設定にsequenceEnabled=trueが含まれている場合にのみ存在します。デフォルトはtrueです。ソートに使用され、生成規則は同期プロセス内で、タイムスタンプ+5桁以下の連番を加算することです。
},
"key": { //UPDATEのみ存在し、変更前の値を表します。
},
"sql": {
"ddl": ""
} //DDLのみ存在します。op = ddlの場合、DDLステートメントが書き込まれます。
}
データ例は以下の通りです:
INSERT(挿入)データの例{ "data":{ "col1":"2020-11-25 00:01:02", "col2":"hello world", "col3":"INTERVAL '3' DAY", "col4":1.2222, "col5":9.999999999999999308, "col6":129, "col7":"00:01:02", "col8":1, "col9":"2020-11-25 00:01:02.012345 Asia/Shanghai", "col10":2147483646, "col11":9223372036854775806, "col12":"aGVsbG8gd29ybGQ=", "col13":"INTERVAL '4' YEAR", "col14":3, "col15":"2020-11-25", "col16":9.9999999999308, "col17":10223372036854775806, "col18":"1606233662.012345" }, "meta":{ "posttime":"2020-12-07T13:22:00", "op":"ins", "size":10, "time":"2020-11-25T00:01:02", "idx":"1/10", "seq":1, "table":"mock_database.mock_table", "rowid":"mock_database.mock_table-3129", "trans":"shareplex_transaction_id", "scn":"123456789" } }UPDATE(更新)データの例{ "data":{ "string":"hello world 2020" }, "meta":{ "posttime":"2020-12-07T13:59:09", "op":"upd", "size":10, "time":"2020-11-25T00:01:02", "idx":"1/10", "seq":1, "table":"mock_database.mock_table", "rowid":"mock_database.mock_table-3\u0001129", "trans":"shareplex_transaction_id", "scn":"123456789" }, "key":{ "col1":"2020-11-25 00:01:02", "col2":"hello world", "col3":"INTERVAL '3' DAY", "col4":1.2222, "col5":9.9999999999999308, "col6":129, "col7":"00:01:02", "col8":1, "col9":"2020-11-25 00:01:02.012345 Asia/Shanghai", "col10":2147483646, "col11":9223372036854775806, "col12":"aGVsbG8gd29ybGQ=", "col13":"INTERVAL '4' YEAR", "col14":3, "col15":"2020-11-25", "col16":9.9999999999308, "col17":10223372036854775806, "col18":"1606233662.012345" } }DELETE(削除)データの例{ "data":{ "col1":"2020-11-25 00:01:02", "col2":"hello world", "col3":"INTERVAL '3' DAY", "col4":1.2222, "col5":9.9999999999308, "col6":129, "col7":"00:01:02", "col8":1, "col9":"2020-11-25 00:01:02.012345 Asia/Shanghai", "col10":2147483646, "col11":9223372036854775806, "col12":"aGVsbG8gd29ybGQ=", "col13":"INTERVAL '4' YEAR", "col14":3, "col15":"2020-11-25", "col16":9.9999999308, "col17":10223372036854775806, "col18":"1606233662.012345" }, "meta":{ "posttime":"2020-12-07T13:34:10", "op":"del", "size":10, "time":"2020-11-25T00:01:02", "idx":"1/10", "seq":1, "table":"mock_database.mock_table", "rowid":"mock_database.mock_table-3\u0001129", "trans":"shareplex_transaction_id", "scn":"123456789" } }DDLの例
ALTER TABLE connector_test.all_mysql_type_test ADD COLUMN c90 VARCHAR(30) DEFAULT "test" COMMENT 'test';{ "data": {}, "meta": { "posttime": "2022-12-16T15:54:51", "op": "ddl", "size": 0, "time": "2022-12-16T15:53:29", "idx": "0/0", "seq": 0, "table": "connector_test.all_mysql_type_test", "rowid": "connector_test.all_mysql_type_test-", "trans": null, "scn": "null" }, "sql": { "ddl": "ALTER TABLE connector_test.all_mysql_type_test ADD COLUMN c90 VARCHAR(30) DEFAULT \"test\" comment 'test'" } }
DefaultExtendColumnType JSONメッセージ形式
データをKafkaに移行する際、DefaultExtendColumnTypeのシリアル化方式では以下のJSONメッセージ形式を使用します。
DefaultExtendColumnType JSONメッセージ形式は、DEFAULT の基礎に、イメージ内にフィールドのデータ型を表す __light_type フィールドを追加します。
{
"prevStruct": { // 変更前イメージ
},
"postStruct": { // 変更後イメージ
"__light_type": {
"col": { // フィールド名
"schemaType": "type" // 値の型
}
}
},
"allMetaData": {
}
}
データ例は以下のとおりです:
INSERT(挿入)データの例{ "allMetaData":{ "checkpoint": null, "record_primary_key": "int8\u0001int16", "source_identity": null, "uniqueId": "{tid:11039xxxx127, partition_id:0, part_cnt:0},5917,391,0", "transId": "{hash:123456, inc:1234, addr:\"1.2.3.4:2883\", t:123456}", "clusterId": "123456", "record_primary_value": "3\u0001129", "dbType": "OB_MYSQL", "table_name": "table", "db": "tenant.database", "timestamp": "1609344671" }, "prevStruct": null, "recordType": "INSERT", "postStruct":{ "col1": 3, "col2": 129, "col3": 2147483646, "col4": 9223372036854775806, "col5": 10223372036854775806, "col6": 1.2222, "col7": 9.99999999999999909326625337248, "col8": "hello world", "col9": "aGVsbG8gd29ybGQ=", "col10": 9.99999999999999909326625337248461995470488734032045693707225049338, "col11": "2020-11-25", "col12": "00:01:02", "col13": "2020-11-25 00:01:02", "col14": "1606233662.012345", "__light_type":{ "int8":{ "schemaType":"TINYINT" }, "int16":{ "schemaType":"SMALLINT" }, "int32":{ "schemaType":"INT" }, "int64":{ "schemaType":"INT64" }, "bigInt":{ "schemaType":"BIGINT" }, "float32":{ "schemaType":"FLOAT" }, "float64":{ "schemaType":"DOUBLE" }, "string":{ "schemaType":"VARCHAR" }, "bytes":{ "schemaType":"BLOB" }, "decimal":{ "schemaType":"DECIMAL" }, "localDate":{ "schemaType":"DATE" }, "localTime":{ "schemaType":"TIME" }, "localDateTime":{ "schemaType":"DATETIME" }, "timestamp_in_long":{ "schemaType":"TIMESTAMP" } } } }recordTypeがROWの場合、全量渡されるデータであり、形式はINSERTと同じです。{ "allMetaData":{ "checkpoint": null, "record_primary_key": "int8\u0001int16", "source_identity": null, "uniqueId": null, "transId": null, "clusterId": null, "record_primary_value": "3\u0001129", "dbType": "OB_MYSQL", "table_name": "table", "db": "tenant.database", "timestamp": null }, "prevStruct": null, "recordType": "ROW", "postStruct":{ "col1": 3, "col2": 129, "col3": 2147483646, "col4": 9223372036854775806, "col5": 10223372036854775806, "col6": 1.2222, "col7": 9.999999, "col8": "hello world", "col9": "aGVsbG8gd29ybGQ=", "col10": 9.99999999999, "col11": "2020-11-25", "col12": "00:01:02", "col13": "2020-11-25 00:01:02", "col14": "1606233662.012345", "__light_type":{ "int8":{ "schemaType":"TINYINT" }, "int16":{ "schemaType":"SMALLINT" }, "int32":{ "schemaType":"INT" }, "int64":{ "schemaType":"INT64" }, "bigInt":{ "schemaType":"BIGINT" }, "float32":{ "schemaType":"FLOAT" }, "float64":{ "schemaType":"DOUBLE" }, "string":{ "schemaType":"VARCHAR" }, "bytes":{ "schemaType":"BLOB" }, "decimal":{ "schemaType":"DECIMAL" }, "localDate":{ "schemaType":"DATE" }, "localTime":{ "schemaType":"TIME" }, "localDateTime":{ "schemaType":"DATETIME" }, "timestamp_in_long":{ "schemaType":"TIMESTAMP" } } } }UPDATE(更新)データの例{ "allMetaData": { "checkpoint": null, "record_primary_key": "int8\u0001int16", "source_identity": null, "uniqueId": "{tid:1103909342127, partition_id:0, part_cnt:0},5917,391,0", "transId": "{hash:123456, inc:1234, addr:\"1.2.3.4:2883\", t:123456}", "clusterId": "123456", "record_primary_value": "3\u0001129", "dbType":"OB_MYSQL", "table_name": "table", "db": "tenant.database", "timestamp": "1609344671" }, "prevStruct": { "col1": 3, "col2": 129, "col3": 2147483646, "col4": 9223372036854775806, "col5": 10223372036854775806, "col6": 1.2222, "col7": 9.999999999999, "col8": "hello world 2020", "col9": "aGVsbG8gd29ybGQ=", "col10": 9.999999999999, "col11": "2020-11-25", "col12": "00:01:02", "col13": "2020-11-25 00:01:02", "col14": "1606233662.012345", "__light_type": { "int8": { "schemaType": "TINYINT" }, "int16": { "schemaType": "SMALLINT" }, "int32": { "schemaType": "INT" }, "int64": { "schemaType": "INT64" }, "bigInt": { "schemaType": "BIGINT" }, "float32": { "schemaType": "FLOAT" }, "float64": { "schemaType": "DOUBLE" }, "string": { "schemaType": "VARCHAR" }, "bytes": { "schemaType": "BLOB" }, "decimal": { "schemaType": "DECIMAL" }, "localDate": { "schemaType": "DATE" }, "localTime": { "schemaType": "TIME" }, "localDateTime": { "schemaType": "DATETIME" }, "timestamp_in_long": { "schemaType": "TIMESTAMP" } } }, "recordType": "UPDATE", "postStruct": { "col1": 3, "col2": 129, "col3": 2147483646, "col4": 9223372036854775806, "col5": 10223372036854775806, "col6": 1.2222, "col7": 9.999999999999, "col8": "hello world 2020", "col9": "aGVsbG8gd29ybGQ=", "col10": 9.999999999999, "col11": "2020-11-25", "col12": "00:01:02", "col13": "2020-11-25 00:01:02", "col14": "1606233662.012345", "__light_type": { "int8": { "schemaType": "TINYINT" }, "int16": { "schemaType": "SMALLINT" }, "int32": { "schemaType": "INT" }, "int64": { "schemaType": "INT64" }, "bigInt": { "schemaType": "BIGINT" }, "float32": { "schemaType": "FLOAT" }, "float64": { "schemaType": "DOUBLE" }, "string": { "schemaType": "VARCHAR" }, "bytes": { "schemaType": "BLOB" }, "decimal": { "schemaType": "DECIMAL" }, "localDate": { "schemaType": "DATE" }, "localTime": { "schemaType": "TIME" }, "localDateTime": { "schemaType": "DATETIME" }, "timestamp_in_long": { "schemaType": "TIMESTAMP" } } } }DELETE(削除)データの例{ "allMetaData":{ "checkpoint": null, "record_primary_key": "int8\u0001int16", "source_identity": null, "uniqueId": "{tid:1103xxxx127, partition_id:0, part_cnt:0},5917,391,0", "transId": "{hash:123456, inc:1234, addr:\"1.2.3.4:2883\", t:123456}", "clusterId": "123456", "record_primary_value": "3\u0001129", "dbType": "OB_MYSQL", "table_name": "table", "db": "tenant.database", "timestamp": "1609344671" }, "prevStruct":{ "col1": 3, "col2": 129, "col3": 2147483646, "col4": 9223372036854775806, "col5": 10223372036854775806, "col6": 1.2222, "col7": 9.999999999999, "col8": "hello world 2020", "col9": "aGVsbG8gd29ybGQ=", "col10": 9.999999999999, "col11": "2020-11-25", "col12": "00:01:02", "col13": "2020-11-25 00:01:02", "col14": "1606233662.012345", "__light_type":{ "int8":{ "schemaType":"TINYINT" }, "int16":{ "schemaType":"SMALLINT" }, "int32":{ "schemaType":"INT" }, "int64":{ "schemaType":"INT64" }, "bigInt":{ "schemaType":"BIGINT" }, "float32":{ "schemaType":"FLOAT" }, "float64":{ "schemaType":"DOUBLE" }, "string":{ "schemaType":"VARCHAR" }, "bytes":{ "schemaType":"BLOB" }, "decimal":{ "schemaType":"DECIMAL" }, "localDate":{ "schemaType":"DATE" }, "localTime":{ "schemaType":"TIME" }, "localDateTime":{ "schemaType":"DATETIME" }, "timestamp_in_long": { "schemaType": "TIMESTAMP" } } }, "recordType":"DELETE", "postStruct":null }DDLの例
ALTER TABLE connector_test.all_mysql_type_test ADD column c90 VARCHAR(30) DEFAULT "test" COMMENT 'test';{ "prevStruct": null, "postStruct": { "ddl": "ALTER TABLE connector_test.all_mysql_type_test ADD column c90 VARCHAR(30) DEFAULT \"test\" COMMENT 'test'", "__light_type": { "ddl": { "schemaType": "VAR_STRING" } } }, "allMetaData": { "checkpoint": "1671177200", "dbType": "OB_MYSQL", "storeDataSequence": null, "db": "connector_test", "timestamp": "1671177209", "uniqueId": null, "ddlType": "ALTER_TABLE", "record_primary_key": null, "source_identity": null, "record_primary_value": null, "table_name": "all_mysql_type_test" }, "recordType": "DDL" }
Debezium JSONメッセージ形式
OceanBaseデータベースのMySQL互換モードのデータをKafkaに移行する際、Debeziumが使用するシリアル化方式は以下のJSONメッセージ形式であり、合計2種類があります。通常、デフォルトではpayloadの構造のみが表示されます。
schemaとpayloadの両方が存在する場合{ "schema": { //payloadフィールド情報を記述する構造体。デフォルトではこの構造体はありません。 "type": "struct", //structはそのフィールド内部にも構造があることを示します。 "optional": false, //このフィールドが含まれている必要があるかどうか。 "fields": [ { "type": "int64", //フィールドの型 "optional": false, //このフィールドが含まれている必要があるかどうか。 "field": "ts_ms" //フィールド名 } ] }, "payload": { "op": "c", //データ変更タイプ。c(フル、挿入)、u(更新)、d(削除)、HEARTBEAT(ハートビートメッセージ)が含まれます。 "source": { "version": "", //データ移行のバージョン。 "connector": "OB_MYSQL", //データソースのタイプ。 "name": "OMS", //固定値 OMS "ts_ms": 0, //データ変更の秒単位のタイムスタンプ。増分データにのみ存在します。 "db": "test", //SQLステートメントを使用して変更を行うデータベースの名前。OceanBaseデータベースの場合は、データベース名のみが存在し、テナント名は不要です。 "table": "testTab", //SQLステートメントを使用して変更を行うテーブルの名前。 "pos": "553132@1668496109" //Binlogファイル内の位置 [Binlogファイル名]@[Binlogファイル名 offset] }, "before": { //変更前のイメージ "column": "value" //キーと値のペア。フルキーと値を含みます。 }, "after": { //変更後のイメージ "column": "value" //キーと値のペア。フルキーと値を含みます。 }, "ts_ms": 1668497367188 //データ処理のタイムスタンプ } }payloadのみが存在する場合{ "payload": { "op": "c", //データ変更タイプ。c(フル、挿入)、u(更新)、d(削除)、HEARTBEAT(ハートビートメッセージ)が含まれます。 "source": { "version": "", //データ移行サービスのバージョン。 "connector": "OB_MYSQL", //データソースのタイプ "name": "OMS", //固定値 OMS "ts_ms": 0, //データ変更の秒単位のタイムスタンプ。増分データにのみ存在する。 "db": "test", //SQLステートメントを使用して変更するデータベースの名前。OceanBaseデータベースの場合は、データベース名のみが存在し、テナント名は不要です。 "table": "testTab", //SQLステートメントを使用して変更するテーブルの名前。 "pos": "553132@16684****" //Binlogファイル内の位置 [Binlogファイル名]@[Binlogファイル名 offset] }, "before": { //変更前のイメージ "column": "value" //キーと値のペア。フルキーと値を含む。 }, "after": { //変更後のイメージ "column": "value" //キーと値のペア。フルキーと値を含む。 }, "ts_ms": 1668497367188 //データ処理のタイムスタンプ } }
データ例は以下の通りです:
INSERT(挿入)データの例{ "schema":{ "optional":false, "type":"STRUCT", "fields":[ { "field":"before", "optional":false, "type":"struct", "fields":[ { "field":"c01", "optional":false, "type":"int32" }, { "field":"c02", "optional":false, "type":"string" }, { "field":"c03", "optional":false, "type":"string" }, { "field":"c04", "optional":false, "type":"bytes" }, { "field":"c05", "optional":false, "type":"int16" }, { "field":"c06", "optional":false, "type":"int16" }, { "field":"c07", "optional":false, "type":"int32" }, { "field":"c08", "optional":false, "type":"int64" }, { "field":"c09", "optional":false, "type":"float64" }, { "field":"c10", "optional":false, "type":"float64" }, { "field":"c11", "optional":false, "type":"string" }, { "field":"c12", "optional":false, "type":"string" }, { "field":"c13", "optional":false, "type":"string" }, { "field":"c14", "optional":false, "type":"string" }, { "field":"c15", "optional":false, "type":"bytes" }, { "field":"c16", "optional":false, "type":"string" }, { "field":"c17", "optional":false, "type":"bytes" }, { "field":"c18", "optional":false, "type":"bytes" }, { "field":"c19", "optional":false, "type":"bytes" }, { "field":"c20", "optional":false, "type":"bytes" }, { "field":"c21", "optional":false, "type":"string" }, { "field":"c22", "optional":false, "type":"int32" }, { "field":"c23", "optional":false, "type":"int64" }, { "field":"c24", "optional":false, "type":"string" }, { "field":"c25", "optional":false, "type":"int32" }, { "field":"c26", "optional":false, "type":"bytes" } ] }, { "field":"after", "optional":false, "type":"struct", "fields":[ { "field":"c01", "optional":false, "type":"int32" }, { "field":"c02", "optional":false, "type":"string" }, { "field":"c03", "optional":false, "type":"string" }, { "field":"c04", "optional":false, "type":"bytes" }, { "field":"c05", "optional":false, "type":"int16" }, { "field":"c06", "optional":false, "type":"int16" }, { "field":"c07", "optional":false, "type":"int32" }, { "field":"c08", "optional":false, "type":"int64" }, { "field":"c09", "optional":false, "type":"float64" }, { "field":"c10", "optional":false, "type":"float64" }, { "field":"c11", "optional":false, "type":"string" }, { "field":"c12", "optional":false, "type":"string" }, { "field":"c13", "optional":false, "type":"string" }, { "field":"c14", "optional":false, "type":"string" }, { "field":"c15", "optional":false, "type":"bytes" }, { "field":"c16", "optional":false, "type":"string" }, { "field":"c17", "optional":false, "type":"bytes" }, { "field":"c18", "optional":false, "type":"bytes" }, { "field":"c19", "optional":false, "type":"bytes" }, { "field":"c20", "optional":false, "type":"bytes" }, { "field":"c21", "optional":false, "type":"string" }, { "field":"c22", "optional":false, "type":"int32" }, { "field":"c23", "optional":false, "type":"int64" }, { "field":"c24", "optional":false, "type":"string" }, { "field":"c25", "optional":false, "type":"int32" }, { "field":"c26", "optional":false, "type":"bytes" } ] }, { "field":"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":"db", "optional":false, "type":"string" }, { "field":"table", "optional":false, "type":"string" }, { "field":"server_id", "optional":false, "type":"int64" }, { "field":"pos", "optional":false, "type":"string" } ] }, { "field":"op", "optional":false, "type":"string" }, { "field":"ts_ms", "optional":false, "type":"int64" } ] }, "payload":{ "op":"c", "source":{ "connector":"OB_MYSQL", "pos":"703223@166849****", "name":"OMS", "version":"", "ts_ms":1668491621000, "db":"test", "table":"table_name" }, "after":{ "c11":"a", "c10":2.4212412, "c13":"c", "c12":"b", "c15":"65", "c14":"d", "c17":"67", "c16":"f", "c19":"69000000000000", "c18":"68", "c20":"6A", "c22":19311, "c21":"2022-11-15T05:12:11Z", "c02":"12312", "c24":1668489131000, "c01":2, "c23":36060000000, "c04":"61", "c26":"6B", "c03":"1241.41000", "c25":2022, "c06":141, "c05":11, "c08":412124124, "c07":4241, "c09":2.11111 }, "ts_ms":1668495423594 } }UPDATE(更新)データの例{ "schema":{ "optional":false, "type":"STRUCT", "fields":[ { "field":"before", "optional":false, "type":"struct", "fields":[ { "field":"c01", "optional":false, "type":"int32" }, { "field":"c02", "optional":false, "type":"string" }, { "field":"c03", "optional":false, "type":"string" }, { "field":"c04", "optional":false, "type":"bytes" }, { "field":"c05", "optional":false, "type":"int16" }, { "field":"c06", "optional":false, "type":"int16" }, { "field":"c07", "optional":false, "type":"int32" }, { "field":"c08", "optional":false, "type":"int64" }, { "field":"c09", "optional":false, "type":"float64" }, { "field":"c10", "optional":false, "type":"float64" }, { "field":"c11", "optional":false, "type":"string" }, { "field":"c12", "optional":false, "type":"string" }, { "field":"c13", "optional":false, "type":"string" }, { "field":"c14", "optional":false, "type":"string" }, { "field":"c15", "optional":false, "type":"bytes" }, { "field":"c16", "optional":false, "type":"string" }, { "field":"c17", "optional":false, "type":"bytes" }, { "field":"c18", "optional":false, "type":"bytes" }, { "field":"c19", "optional":false, "type":"bytes" }, { "field":"c20", "optional":false, "type":"bytes" }, { "field":"c21", "optional":false, "type":"string" }, { "field":"c22", "optional":false, "type":"int32" }, { "field":"c23", "optional":false, "type":"int64" }, { "field":"c24", "optional":false, "type":"string" }, { "field":"c25", "optional":false, "type":"int32" }, { "field":"c26", "optional":false, "type":"bytes" } ] }, { "field":"after", "optional":false, "type":"struct", "fields":[ { "field":"c01", "optional":false, "type":"int32" }, { "field":"c02", "optional":false, "type":"string" }, { "field":"c03", "optional":false, "type":"string" }, { "field":"c04", "optional":false, "type":"bytes" }, { "field":"c05", "optional":false, "type":"int16" }, { "field":"c06", "optional":false, "type":"int16" }, { "field":"c07", "optional":false, "type":"int32" }, { "field":"c08", "optional":false, "type":"int64" }, { "field":"c09", "optional":false, "type":"float64" }, { "field":"c10", "optional":false, "type":"float64" }, { "field":"c11", "optional":false, "type":"string" }, { "field":"c12", "optional":false, "type":"string" }, { "field":"c13", "optional":false, "type":"string" }, { "field":"c14", "optional":false, "type":"string" }, { "field":"c15", "optional":false, "type":"bytes" }, { "field":"c16", "optional":false, "type":"string" }, { "field":"c17", "optional":false, "type":"bytes" }, { "field":"c18", "optional":false, "type":"bytes" }, { "field":"c19", "optional":false, "type":"bytes" }, { "field":"c20", "optional":false, "type":"bytes" }, { "field":"c21", "optional":false, "type":"string" }, { "field":"c22", "optional":false, "type":"int32" }, { "field":"c23", "optional":false, "type":"int64" }, { "field":"c24", "optional":false, "type":"string" }, { "field":"c25", "optional":false, "type":"int32" }, { "field":"c26", "optional":false, "type":"bytes" } ] }, { "field":"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":"db", "optional":false, "type":"string" }, { "field":"table", "optional":false, "type":"string" }, { "field":"server_id", "optional":false, "type":"int64" }, { "field":"pos", "optional":false, "type":"string" } ] }, { "field":"op", "optional":false, "type":"string" }, { "field":"ts_ms", "optional":false, "type":"int64" } ] }, "payload":{ "op":"u", "before":{ "c11":"a", "c10":2.4212412, "c13":"c", "c12":"b", "c15":"65", "c14":"d", "c17":"67", "c16":"f", "c19":"6900000000", "c18":"68", "c20":"6A", "c22":19311, "c21":"2022-11-15T05:12:11Z", "c02":"12312", "c24":1668489131000, "c01":1, "c23":36060000000, "c04":"61", "c26":"6B", "c03":"1241.41000", "c25":2022, "c06":141, "c05":11, "c08":412124124, "c07":4241, "c09":2.11111 }, "source":{ "connector":"OB_MYSQL", "pos":"436999@166849****", "name":"OMS", "version":"", "ts_ms":1668495861000, "db":"test", "table":"table_name" }, "after":{ "c11":"aa", "c10":2.4212412, "c13":"c", "c12":"b", "c15":"65", "c14":"d", "c17":"67", "c16":"f", "c19":"69000000000", "c18":"68", "c20":"6A", "c22":19311, "c21":"2022-11-15T05:12:11Z", "c02":"12312", "c24":1668489131000, "c01":1, "c23":36060000000, "c04":"61", "c26":"6B", "c03":"1241.41000", "c25":2022, "c06":141, "c05":11, "c08":412124124, "c07":4241, "c09":2.11111 }, "ts_ms":1668495906356 } }DELETE(削除)データの例{ "schema":{ "optional":false, "type":"STRUCT", "fields":[ { "field":"before", "optional":false, "type":"struct", "fields":[ { "field":"c01", "optional":false, "type":"int32" }, { "field":"c02", "optional":false, "type":"string" }, { "field":"c03", "optional":false, "type":"string" }, { "field":"c04", "optional":false, "type":"bytes" }, { "field":"c05", "optional":false, "type":"int16" }, { "field":"c06", "optional":false, "type":"int16" }, { "field":"c07", "optional":false, "type":"int32" }, { "field":"c08", "optional":false, "type":"int64" }, { "field":"c09", "optional":false, "type":"float64" }, { "field":"c10", "optional":false, "type":"float64" }, { "field":"c11", "optional":false, "type":"string" }, { "field":"c12", "optional":false, "type":"string" }, { "field":"c13", "optional":false, "type":"string" }, { "field":"c14", "optional":false, "type":"string" }, { "field":"c15", "optional":false, "type":"bytes" }, { "field":"c16", "optional":false, "type":"string" }, { "field":"c17", "optional":false, "type":"bytes" }, { "field":"c18", "optional":false, "type":"bytes" }, { "field":"c19", "optional":false, "type":"bytes" }, { "field":"c20", "optional":false, "type":"bytes" }, { "field":"c21", "optional":false, "type":"string" }, { "field":"c22", "optional":false, "type":"int32" }, { "field":"c23", "optional":false, "type":"int64" }, { "field":"c24", "optional":false, "type":"string" }, { "field":"c25", "optional":false, "type":"int32" }, { "field":"c26", "optional":false, "type":"bytes" } ] }, { "field":"after", "optional":false, "type":"struct", "fields":[ { "field":"c01", "optional":false, "type":"int32" }, { "field":"c02", "optional":false, "type":"string" }, { "field":"c03", "optional":false, "type":"string" }, { "field":"c04", "optional":false, "type":"bytes" }, { "field":"c05", "optional":false, "type":"int16" }, { "field":"c06", "optional":false, "type":"int16" }, { "field":"c07", "optional":false, "type":"int32" }, { "field":"c08", "optional":false, "type":"int64" }, { "field":"c09", "optional":false, "type":"float64" }, { "field":"c10", "optional":false, "type":"float64" }, { "field":"c11", "optional":false, "type":"string" }, { "field":"c12", "optional":false, "type":"string" }, { "field":"c13", "optional":false, "type":"string" }, { "field":"c14", "optional":false, "type":"string" }, { "field":"c15", "optional":false, "type":"bytes" }, { "field":"c16", "optional":false, "type":"string" }, { "field":"c17", "optional":false, "type":"bytes" }, { "field":"c18", "optional":false, "type":"bytes" }, { "field":"c19", "optional":false, "type":"bytes" }, { "field":"c20", "optional":false, "type":"bytes" }, { "field":"c21", "optional":false, "type":"string" }, { "field":"c22", "optional":false, "type":"int32" }, { "field":"c23", "optional":false, "type":"int64" }, { "field":"c24", "optional":false, "type":"string" }, { "field":"c25", "optional":false, "type":"int32" }, { "field":"c26", "optional":false, "type":"bytes" } ] }, { "field":"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":"db", "optional":false, "type":"string" }, { "field":"table", "optional":false, "type":"string" }, { "field":"server_id", "optional":false, "type":"int64" }, { "field":"pos", "optional":false, "type":"string" } ] }, { "field":"op", "optional":false, "type":"string" }, { "field":"ts_ms", "optional":false, "type":"int64" } ] }, "payload":{ "op":"d", "before":{ "c11":"aa", "c10":2.4212412, "c13":"c", "c12":"b", "c15":"65", "c14":"d", "c17":"67", "c16":"f", "c19":"69000000000", "c18":"68", "c20":"6A", "c22":19311, "c21":"2022-11-15T05:12:11Z", "c02":"12312", "c24":1668489131000, "c01":1, "c23":36060000000, "c04":"61", "c26":"6B", "c03":"1241.41000", "c25":2022, "c06":141, "c05":11, "c08":412124124, "c07":4241, "c09":2.11111 }, "source":{ "connector":"OB_MYSQL", "pos":"553132@1668****", "name":"OMS", "version":"", "ts_ms":1668496109000, "db":"test", "table":"table_name" }, "ts_ms":1668496119717 } }
DebeziumFlatten JSONメッセージ形式
OceanBaseデータベースのMySQL互換モードのデータをKafkaに移行する際、DebeziumFlatten形式でシリアライズされるJSONメッセージの形式は以下のとおりです。Debezium形式と比較して、schemaとpayloadが埋め込まれなくなっています。
{
"op": "c", //データ変更タイプ。c(フル、挿入)、u(更新)、d(削除)、HEARTBEAT(ハートビートメッセージ)が含まれます。
"source": {
"version": "", //データ移行サービスのバージョン。
"connector": "OB_MYSQL", //データソースのタイプ。
"name": "OMS", //固定値 OMS
"ts_ms": 0, //データ変更の秒単位のタイムスタンプ。増分データにのみ存在する。
"db": "test", //SQLステートメントを使用して変更するデータベースの名前。OceanBaseデータベースの場合は、データベース名のみが存在し、テナント名は不要です。
"table": "testTab", //SQLステートメントを使用して変更するテーブルの名前。
"pos": "553132@16684****" //binlogファイル内の位置 [binlogファイル名]@[binlogファイル名 offset]
},
"before": { //変更前のイメージ
"column": "value" //キーと値のペア。フルキーと値を含む。
},
"after": { // 変更後のイメージ
"column": "value" // キーと値のペア。フルキー・フルバリューを含む
},
"ts_ms": 1668497367188 // データ処理タイムスタンプ
}
データ例は以下の通りです:
INSERT(挿入)データの例{ "op":"c", "source":{ "connector":"OB_MYSQL", "pos":"703223@166849****", "name":"OMS", "version":"", "ts_ms":1668491621000, "db":"test", "table":"table_name" }, "after":{ "c11":"a", "c10":2.4212412, "c13":"c", "c12":"b", "c15":"65", "c14":"d", "c17":"67", "c16":"f", "c19":"69000000000000000000", "c18":"68", "c20":"6A", "c22":19311, "c21":"2022-11-15T05:12:11Z", "c02":"12312", "c24":1668489131000, "c01":2, "c23":36060000000, "c04":"61", "c26":"6B", "c03":"1241.41000", "c25":2022, "c06":141, "c05":11, "c08":412124124, "c07":4241, "c09":2.11111 }, "ts_ms":1668495423594 }UPDATE(更新)データの例{ "op":"u", "before":{ "c11":"a", "c10":2.4212412, "c13":"c", "c12":"b", "c15":"65", "c14":"d", "c17":"67", "c16":"f", "c19":"690000000000000000000000000000000", "c18":"68", "c20":"6A", "c22":19311, "c21":"2022-11-15T05:12:11Z", "c02":"12312", "c24":1668489131000, "c01":1, "c23":36060000000, "c04":"61", "c26":"6B", "c03":"1241.41000", "c25":2022, "c06":141, "c05":11, "c08":412124124, "c07":4241, "c09":2.11111 }, "source":{ "connector":"OB_MYSQL", "pos":"436999@166849****", "name":"OMS", "version":"", "ts_ms":1668495861000, "db":"test", "table":"table_name" }, "after":{ "c11":"aa", "c10":2.4212412, "c13":"c", "c12":"b", "c15":"65", "c14":"d", "c17":"67", "c16":"f", "c19":"69000000000000000000000000", "c18":"68", "c20":"6A", "c22":19311, "c21":"2022-11-15T05:12:11Z", "c02":"12312", "c24":1668489131000, "c01":1, "c23":36060000000, "c04":"61", "c26":"6B", "c03":"1241.41000", "c25":2022, "c06":141, "c05":11, "c08":412124124, "c07":4241, "c09":2.11111 }, "ts_ms":1668495906356 }DELETE(削除)データの例{ "op":"d", "before":{ "c11":"aa", "c10":2.4212412, "c13":"c", "c12":"b", "c15":"65", "c14":"d", "c17":"67", "c16":"f", "c19":"69000000000000000000000000000", "c18":"68", "c20":"6A", "c22":19311, "c21":"2022-11-15T05:12:11Z", "c02":"12312", "c24":1668489131000, "c01":1, "c23":36060000000, "c04":"61", "c26":"6B", "c03":"1241.41000", "c25":2022, "c06":141, "c05":11, "c08":412124124, "c07":4241, "c09":2.11111 }, "source":{ "connector":"OB_MYSQL", "pos":"553132@1668****", "name":"OMS", "version":"", "ts_ms":1668496109000, "db":"test", "table":"table_name" }, "ts_ms":1668496119717 }
DebeziumSmt JSONメッセージ形式
DebeziumSmtは、Debeziumが提供する設定方法の一つで、イベントのフラット化されたシングルメッセージ変換(Single Message Transform、SMT)を用いて単一の情報を変換・処理します。OceanBaseデータベースのMySQL互換モードのデータをKafkaに移行する際、DebeziumSmtのシリアル化方式では、JSONメッセージ形式にafterのkey:valueのみが表示されます。
例えば、DebeziumSmtシリアル化方式を使用してデータを更新する場合:
{
"op": "u",
"source": {
"connector": "OB_MYSQL",
"name": "OMS"
},
"ts_ms": 1668496119717,
"before": {
"field1": "before_value1",
"field2": "before_value2"
},
"after": {
"field1": "after_value1",
"field2": "after_value2"
}
}
SMTが上記の例のメッセージを処理した後、メッセージ形式は簡略化されます。DebeziumSmtシリアル化方式を使用する場合、JSONメッセージ形式は以下のようになります。
{
"field1": "after_value1",
"field2": "after_value2"
}
データ例は以下の通りです:
INSERT(挿入)データの例{ "field1": "after_value1", "field2": "after_value2", "__deleted": "false" }UPDATE(更新)データの例{ "field1": "after_value1", "field2": "after_value2", "__deleted": "false" }DELETE(削除)データの例{ "field1": "after_value1", "field2": "after_value2", "__deleted": "true" }
Avro JSONメッセージ形式
OceanBaseデータベースのMySQL互換モードのデータをKafkaに移行する際、シリアル化方式としてAvroを使用する場合のJSONメッセージ形式は以下のとおりです。
フル移行
{ "version": 1, "id": 0, "sourceTimestamp": 1702371565, // タイムスタンプのセーフポイント "sourcePosition": "", // フル移行ではpositionなどの情報はありません "safeSourcePosition": "", "sourceTxid": "", "source": { "sourceType": "MySQL", // 固定値 MySQL "version": "OBMySQL" // 固定値 OBMySQL }, "operation": "INIT", // フル移行のタイプはINITです "objectName": "test***", "processTimestamps": [ 1702371565238 ], // 配信時間のみ "tags": { "pk_uk_info": "{\"PRIMARY\":[\"id\"]}" // 主キーのみのタイプ }, "fields": [ { "name": "id", "dataTypeNumber": 246 }, // 各列の型 { "name": "bid", "dataTypeNumber": 3 }, { "name": "name", "dataTypeNumber": 15 }, { "name": "address", "dataTypeNumber": 254 } ], "beforeImages": null, // フル移行前のイメージは空です "afterImages": [ // 移行後のイメージ。INTEGER型のprecisionは8、FLOAT型のprecisionは8、scaleは64です { "value": "1", "precision": 1, "scale": 0 }, { "precision": 8, "value": "11" }, { "charset": "utf8mb4", "value": { "bytes": "yyy" } }, null ] }増分同期DML
INSERT(挿入)データの例{ "version": 1, "id": 170236922143600000, "sourceTimestamp": 1702369092, "sourcePosition": "1702369080", // OceanBaseデータベースMySQL互換モードのチェックポイント "safeSourcePosition": "1702369080", // OceanBaseデータベースMySQL互換モードのチェックポイント "sourceTxid": "", "source": { "sourceType": "MySQL", "version": "OBMySQL" }, "operation": "INSERT", "objectName": "test***", "processTimestamps": [1702369221480], "tags": { "pk_uk_info": "{\"PRIMARY\":[\"id\"]}" }, "fields": [ {"name": "id", "dataTypeNumber": 8}, {"name": "bid", "dataTypeNumber": 3}, {"name": "name", "dataTypeNumber": 15} ], "beforeImages": null, // INSERT前のイメージは空 "afterImages": [ {"precision": 8, "value": "2"}, {"precision": 8, "value": "12"}, {"charset": "utf8mb4", "value": {"bytes": "xxx"} } ] }データを
UPDATEする例{ "version": 1, "id": 170236975822100001, "sourceTimestamp": 1702369757, "sourcePosition": "1702369756", "safeSourcePosition": "1702369756", "sourceTxid": "", "source": { "sourceType": "MySQL", "version": "OBMySQL" }, "operation": "UPDATE", "objectName": "test***", "processTimestamps": [1702369758237], "tags": { "pk_uk_info": "{\"PRIMARY\":[\"id\"]}" }, "fields": [ {"name": "id", "dataTypeNumber": 8}, {"name": "bid", "dataTypeNumber": 3}, {"name": "name", "dataTypeNumber": 15} ], "beforeImages": [ // UPDATEの前後のイメージ {"precision": 8, "value": "3"}, {"precision": 8, "value": "22"}, {"charset": "utf8mb4", "value": {"bytes": "xxx"}} ], "afterImages": [ {"precision": 8, "value": "3"}, {"precision": 8, "value": "44"}, {"charset": "utf8mb4", "value": {"bytes": "xxx"}} ] }データを
DELETEする例{ "version": 1, "id": 170236976527500000, "sourceTimestamp": 1702369764, "sourcePosition": "1702369763", "safeSourcePosition": "1702369763", "sourceTxid": "", "source": { "sourceType": "MySQL", "version": "OBMySQL" }, "operation": "DELETE", "objectName": "test***", "processTimestamps": [1702369765287], "tags": { "pk_uk_info": "{\"PRIMARY\":[\"id\"]}" }, "fields": [ {"name": "id", "dataTypeNumber": 8}, {"name": "bid", "dataTypeNumber": 3}, {"name": "name", "dataTypeNumber": 15} ], "beforeImages": [ {"precision": 8, "value": "3"}, {"precision": 8, "value": "44"}, {"charset": "utf8mb4", "value": {"bytes": "xxx"}} ], "afterImages": null // DELETE後のイメージは空 }
増分同期DDL
{ "version": 1, "id": 170236979372400000, "sourceTimestamp": 1702369793, "sourcePosition": "1702369792", "safeSourcePosition": "1702369792", "sourceTxid": "", "source": { "sourceType": "MySQL", "version": "OBMySQL" }, "operation": "DDL", "objectName": "test***", "processTimestamps": [1702369794543], "tags": {}, "fields": null, // 増分同期DDLにはfieldsとbeforeImagesがありません "beforeImages": null, "afterImages": "alter table multi_db_multi_tbl add column address char(20) default null" // STRING型のafterImagesはDDLステートメントです }
データベースからテキストプロトコルへの形式説明
OceanBaseデータベースのデータをKafkaに移行する場合:
シリアライズ方式が Default、Canal、DataWorks(V2.0以降をサポート)、SharePlex、または DefaultExtendColumnType の場合、OceanBaseデータベースの2つの互換モードに対応するマッピング説明は以下のとおりです。
OceanBaseデータベース MySQL互換モード
データ型マッピングタイプ説明TINYINT
SMALLINT
MEDIUMINT
INT
INTEGER
YEAR
BOOL
BOOLEANLong 64ビット以下の整数型。
通常の数値、例えば1000は、科学記数法を使用しません。
BOOL/BOOLEANの場合、true = 1、false = 0。DECIMAL
NUMERICBigDecimal 厳密小数値型および64ビットを超える整数型。整数値には小数点や小数部が表示されません。小数を含む値については、データベースから渡されたデータに基づいて桁数が表示され、末尾の0は切り捨てず、科学記数法が使用されます。 FLOAT
DOUBLEDouble 浮動小数点数
ソース側がFLOATまたはDOUBLE型かどうかによって有効数字の位数が決定します。FLOATは7桁、DOUBLEは16桁の有効数字です。CHAR
VARCHAR
TINYTEXT
TEXT
MEDIUMTEXT
LONGTEXT
ENUM
SETString 文字列。 TINYBLOB
BLOB
MEDIUMBLOB
LONGBLOB
BINARY
VARBINARY
BITBytes バイト配列で、デフォルトではBASE64エンコードで表示されます(CanalはデフォルトでLATIN1エンコードで表示され、BIT型は数値に変換されます)。 説明
BIT定長型の場合、増分でバイト配列を受信すると上位の0が削除されますが、フル量では削除されないため、表示されるBASE64エンコードが一致しない場合があります。しかし、実際の結果は一致しており、デコード後の結果も一致します。
DATE Date 日付型で、形式は YYYY-MM-DDです。無効な時間の場合、元の文字列が表示されます。TIME Time 時刻型で、形式は HH:mm:ss[.nnnnnnnnn]です。
秒未満の時間は最大9桁まで表示されます。秒未満の時間の場合、すべての非0の数値が表示されます。無効な時間の場合、元の文字列が表示されます。DATETIME DateTime 日時型で、タイムゾーンを含みます。形式は YYYY-MM-DD HH:mm:ss[.nnnnnnnnn] [zoneId]です。
秒未満の時間は最大9桁まで表示されます。秒未満の時間の場合、すべての非0の数値が表示されます。無効な時間の場合、元の文字列が表示されます。TIMESTAMP Timestamp タイムスタンプ型で、形式は [秒単位のタイムスタンプ][.nnnnnnnnn]です(Canalの形式はYYYY-MM-DD HH:mm:ss[.nnnnnnnnn])。
秒未満の時間は最大9桁まで表示されます。秒未満の時間の場合、すべての非0の数値が表示されます。無効な時間の場合、0000-00-00 00:00:00の形式で表示されます。OceanBaseデータベース Oracle互換モード
データ型マッピングタイプ説明INTEGER Long 64ビット以下の整数型。
通常の数値、例えば1000は、科学記数法を使用しません。NUMBER
FLOATBigDecimal 厳密な小数値型および64ビットを超える整数型。 BINARY_FLOAT BINARY_DOUBLE Double 浮動小数点数
ソースがFLOAT型かDOUBLE型かに応じて有効数字の位数が決まります。FLOATは7桁、DOUBLEは16桁の有効数字です。VARCHAR2
NVARCHAR2
INTERVAL YEAR TO MONTH
INTERVAL DAY TO SECOND
CLOB
NCLOB
ROWID
UROWIDString 文字列 BLOB
BFILE
RAWBytes バイト配列
デフォルトではBASE64エンコードで表示されます。DATE
TIMESTAMP
TIMESTAMP WITH TIME ZONE
TIMESTAMP WITH LOCAL TIME ZONEDateTime 日付時刻型で、タイムゾーンを含みます。形式は YYYY-MM-DD HH:mm:ss[.nnnnnnnnn] [zoneId]です。
秒未満の時間は最大9桁まで表示されます。秒未満の時間の場合、すべての非0の数字が表示されます。無効な時間の場合、元の文字列が表示されます。
シリアライズ方式が Debezium の場合、OceanBaseデータベースのMySQL互換モードに対応するマッピングの説明は以下のとおりです。
注意
OceanBaseデータベースのOracle互換モードのデータをKafkaに移行する際、シリアライズ方式 Debezium を選択することはサポートされていません。
データ型映射タイプ説明BOOLEAN
BOOLBOOLEAN 値は true と false を含みます。 TINYINT
SMALLINT
MEDIUMINT
INT/INTEGER
BIGINT
YEARLONG -2^63^ ~ 2^63^の範囲の整数型です。 BIGINT STRING 文字列を使用してデータを完全に表示します。 FLOAT
DOUBLEDOUBLE 浮動小数点数。 DECIMAL
NUMERICSTRING 文字列を使用してデータを完全に表示します。小数が含まれる数値については、データベースから渡されたデータに基づいて桁数を表示し、末尾の0を切り捨てずに科学記数法を使用します。 BIT
BINARY
VARBINARY
TINYBLOB
BLOB
MEDIUMBLOB
LONGBLOBBYTES バイト配列、base16エンコード。 CHAR
VARCHAR
TINYTEXT
TEXT
MEDIUMTEXT
LONGTEXT
ENUM
SETSTRING 文字列。 TIMESTAMP STRING 形式はYYYY-MM-DDTHH:mm:ss[.nnnnnnnnn]Z、タイムゾーンは0タイムゾーンです。 DATE LONG 1970-01-01以降の日数を表します。 TIME LONG 00:00:00以降の時間値(マイクロ秒単位)を表します。タイムゾーン情報は含まれません。 DATETIME LONG 1970-01-01 00:00:00以降のミリ秒数を表します。タイムゾーン情報は含まれません。 シリアライズ方式が Avro の場合、OceanBaseデータベースのMySQL互換モードに対応するマッピングの説明は以下のとおりです。
注意
OceanBaseデータベースのMySQL互換モードのデータをKafkaに移行する場合にのみ、シリアライズ方式 Avro を選択できます。
クラス名マッピングタイプTINYINT
BOOLEAN
SMALLINT
MEDIUMINT
INT
BIGINT
BITINTEGER FLOAT
DOUBLEFLOAT DECIMAL
NUMERICDECIMAL VARCHAR
CHAR
TINYTEXT
MEDIUMTEXT
LONGTEXT
TEXTCHARACTER BINARY
VARBINARY
TINYBLOB
MEDIUMBLOB
LONGBLOB
BLOBBinaryObject TIMESTAMP TimestampObject
説明 TIMESTAMP型については、フル移行と増分同期の両方でタイムスタンプに変換されます。無効な時間は-9223372022400Lです。
無効な時間を除き、JavaのInstant.ofEpochSecond(ts, nanos)メソッドを使用して正しい現在時刻を取得できます。DATE
TIME
DATETIME
YEARDATETIME JSON
ENUM
SETTextObject GEOMETRY TextGeometry
説明 現在、データ移行ではEWKT形式がそのまま渡されるため、TextGeometry型にマッピングされます。