本記事では、シリアライズ方式およびデータベースからテキストプロトコルへのデータ形式について説明します。
シリアライズ方式の形式説明
データ移行サービスを使用してデータ送信元のデータを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 のレコードでは、「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", //Schema名。Schemaが存在するシステムで記載必須
"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互換モードの checkpoint "safeSourcePosition": "1702369080", // OceanBaseデータベースMySQL互換モードの checkpoint "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、タイムゾーンはUTC+0です。 DATE LONG 1970年1月1日以降の日数を表します。 TIME LONG 00:00:00以降の時間値(マイクロ秒単位)を表します。タイムゾーン情報は含まれません。 DATETIME LONG 1970年1月1日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タイプにマッピングされます。