Apendiks ini menjelaskan operasi yang didukung, strategi sharding, format data, serta contoh pesan untuk berbagai tipe data DataHub.
Operasi yang didukung untuk berbagai tipe data
Topik merupakan unit dasar untuk langganan dan publikasi data di DataHub, yang merepresentasikan kumpulan data streaming. DataHub mendukung dua tipe data: TUPLE dan BLOB.
|
DataHub type |
Write DML messages |
Write upstream heartbeat messages |
Write DDL messages |
Source to topic mapping |
Type |
|
TUPLE |
Didukung |
Tidak didukung |
Tidak didukung |
Satu tabel ke satu topik |
Tipe yang didukung DataHub |
|
BLOB |
Didukung |
Didukung |
Didukung |
Satu database (beberapa tabel) ke satu topik |
Data biner BLOB |
-
Topik TUPLE memiliki skema tetap yang tidak dapat diubah setelah dibuat. Topik ini cocok untuk skenario di mana skema tabel sumber stabil dan tidak melibatkan operasi DDL seperti
ADD COLUMNatauDROP COLUMN. Topik TUPLE tidak mendukung penerusan pesan DDL atau pesan heartbeat dari sumber ke konsumen downstream. Pemetaan satu-ke-satu ini dapat menyulitkan konsumsi downstream jika Anda memiliki banyak tabel sumber, karena Anda harus membuat topik terpisah untuk masing-masing tabel. -
Topik BLOB tidak memiliki skema yang telah ditentukan sebelumnya dan hanya menyimpan data biner mentah, sehingga lebih fleksibel. Topik ini mendukung penerusan pesan DDL dan pesan heartbeat dari sumber ke konsumen downstream. Topik ini memetakan beberapa tabel dari satu database ke satu topik, memungkinkan Anda menggunakan satu topik untuk semua tabel sumber dan menyederhanakan konsumsi downstream. Tipe ini ideal untuk skenario di mana DataHub berperan sebagai antrian pesan perantara dalam migrasi seluruh database.
Strategi sharding untuk berbagai tipe data
Shard adalah saluran konkuren untuk transmisi data dalam suatu topik DataHub. Meskipun satu shard memiliki throughput tulis yang terbatas, Anda dapat menggunakan beberapa shard untuk meningkatkannya. DataHub hanya menjamin konsumsi berurutan dalam satu shard, bukan di antara beberapa shard. Untuk meningkatkan performa penulisan dengan beberapa shard sekaligus mempertahankan urutan pesan dan mencegah kesenjangan data, DataHub menyediakan strategi sharding berikut untuk tipe data TUPLE dan BLOB.
|
Scenario |
TUPLE |
BLOB |
|
Dengan primary key (termasuk custom primary key) |
Shard berdasarkan primary key |
Shard berdasarkan primary key |
|
Jaminan pengurutan |
Pesan dengan primary key yang sama diproses secara berurutan |
Pesan dengan primary key yang sama diproses secara berurutan |
|
Tanpa primary key |
Random sharding |
Shard berdasarkan nama tabel |
|
Jaminan pengurutan |
Pengurutan tidak dijamin |
Pesan dari tabel yang sama diproses secara berurutan |
Format data
-
TUPLE
Format TUPLE menggunakan tipe data yang secara native didukung oleh DataHub. Data Integration menambahkan beberapa kolom metadata saat Anda membuat topik. Skema untuk format TUPLE mencakup kolom metadata dan kolom data bisnis. Bidang yang diawali garis bawah, seperti
_sequence_id_dan_operation_type_, merupakan kolom metadata. Bidang lainnya adalah kolom data bisnis. Kolom metadata mencakup_sequence_id_,_excute_time_,_source_table_,_before_image_, dan_after_image_.Parameter
Description
_sequence_id_
ID pesan unik bertipe STRING yang terdiri dari angka. Operasi UPDATE_BEFOR dan UPDATE_AFTER untuk pembaruan yang sama berbagi ID sequence yang sama.
_excute_time_
Waktu ketika data dihasilkan.
_source_table_
Nama tabel sumber.
_before_image_
Pre-image. Nilainya
Yuntuk operasi UPDATE_BEFOR atau DELETE, danNuntuk operasi UPDATE_AFTER atau INSERT._after_image_
Post-image. Nilainya
Nuntuk operasi UPDATE_BEFOR atau DELETE, danYuntuk operasi UPDATE_AFTER atau INSERT.Contoh: Tabel berikut menunjukkan data yang disinkronkan ke DataHub setelah menjalankan pernyataan INSERT, UPDATE, dan DELETE.
_sequence_id_
_operation_type_
_excute_time_
_before_image_
_after_image_
1649991610688000000
I
1649991726000
N
Y
1649991610688000001
U
1649991756000
Y
N
1649991610688000001
U
1649991756000
N
Y
1649991610688000002
D
1649991774000
Y
N
-
BLOB
Pesan BLOB adalah data biner yang dibuat dengan mengonversi string JSON. Format JSON yang sesuai adalah sebagai berikut:
{ "schema": { // Metadata tentang perubahan, hanya menentukan nama dan tipe kolom. "dataColumn": [ // Informasi tentang kolom data yang berubah, digunakan untuk memperbarui catatan di tabel tujuan. { "name": "id", "type": "LONG" }, { "name": "name", "type": "STRING" }, { "name": "binData", "type": "BYTES" }, { "name": "ts", "type": "DATE" } ], "primaryKey": [ "pkName1", "pkName2" ], "source": { "dbType": "mysql", "dbVersion": "1.0.0", "dbName": "myDatabase", "schemaName": "mySchema", "tableName": "tableName" } }, "payload": { "before": { "dataColumn":{ "id": 111, "name":"scooter", "binData": "[base64 string]", "ts": 1590315269000 } }, "after": { "dataColumn":{ "id": 222, "name":"donald", "binData": "[base64 string]", "ts": 1590315269000 } }, "sequenceId":XXX, // String yang digunakan untuk pengurutan data saat menggabungkan data penuh dan inkremental. "op": "INSERT/UPDATE/DELETE/TRANSACTION_BEGIN/TRANSACTION_END/CREATE/ALTER/ERASE/QUERY/TRUNCATE/RENAME/CINDEX/DINDEX/GTID/XACOMMIT/XAROLLBACK/MHEARTBEAT...", // Case-sensitive. "timestamp": { "eventTime": 1, // Wajib. Waktu perubahan catatan. Timestamp 13 digit dengan presisi milidetik. "systemTime": 2, // Opsional. Ada untuk beberapa sumber data seperti Oracle CDC. "checkpointTime": 3 // Opsional. Termasuk untuk beberapa sumber data seperti OceanBase. }, "ddl": { "text": "ADD COLUMN ...", "ddlMeta": "[SQLStatement serialized binary, expressed in base64 string]" } }, "version":"1.0.0" }-
Bidang BLOB
PentingTipe data untuk semua bidang dalam pesan didefinisikan oleh StreamX dan mencakup
BOOLEAN,DOUBLE,DATE,BYTES,LONG, danSTRING.BOOLEAN: Nilainya `true` atau `false`. DATE: Nilainya integer 13 digit yang merepresentasikan timestamp dengan presisi milidetik. BYTES: Menyimpan array byte sebagai string yang diencode Base64. Gunakan API `java.util.Base64` untuk encoding dan decoding Base64: String text = "测试text123"; // Encode Base64.getEncoder().encodeToString(text.getBytes("UTF-8")) // Decode Base64.getDecoder().decode(encodedText)Elemen tingkat atas
Elemen tingkat kedua
Description
schema
dataColumn
JSONArray yang berisi informasi tipe untuk kolom data. dataColumn mencatat semua kolom dan tipe datanya dalam catatan perubahan data upstream. Operasi perubahan dapat berupa modifikasi data (seperti insert, delete, atau update) atau modifikasi struktur tabel.
-
name: Nama kolom.
-
type: Tipe data kolom.
primaryKey
Daftar string yang merepresentasikan nama kolom kunci primer.
pk: Nama kunci primer.
source
Objek yang berisi informasi tentang database atau tabel sumber.
-
dbType: String yang merepresentasikan tipe database.
-
dbVersion: String yang merepresentasikan versi database.
-
dbName: String yang merepresentasikan nama database.
-
schemaName: String yang merepresentasikan nama skema, wajib untuk database seperti PostgreSQL dan SQL Server.
-
tableName: String yang merepresentasikan nama tabel.
payload
before
JSONObject yang berisi pre-image data. Untuk operasi
UPDATEpada sumber MySQL, bidangbeforemenyimpan konten catatan sebelum pembaruan.-
Bidang ini diisi saat pesan pembaruan atau penghapusan dibaca dari sumber.
-
dataColumn: Parameter bertipe JSONObject yang merepresentasikan informasi data. Formatnya adalah nama kolom: nilai kolom. Nama kolom adalah string, dan nilai kolom bergantung pada tipe datanya: nilai tipe BYTES direpresentasikan sebagai string Base64, nilai tipe DATE direpresentasikan sebagai timestamp 13 digit bertipe long, dan nilai tipe lain direpresentasikan dengan tipe nativenya.
after
Post-image data. Formatnya sama dengan bidang
before.CatatanBidang ini wajib untuk operasi
UPDATEdanINSERT.op
Tipe operasi. Nilai yang valid:
-
INSERT: Penyisipan data.
-
UPDATE_BEFOR: Pre-image pembaruan.
-
UPDATE_AFTER: Post-image pembaruan.
-
DELETE: Penghapusan data.
-
TRANSACTION_BEGIN: Awal transaksi database.
-
TRANSACTION_END: Akhir transaksi database.
-
CREATE: Pembuatan tabel.
-
ALTER: Perubahan struktur tabel.
-
QUERY: SQL asli untuk perubahan database.
-
TRUNCATE: Mengosongkan tabel.
-
RENAME: Penggantian nama tabel.
-
CINDEX: Pembuatan indeks.
-
DINDEX: Penghapusan indeks.
-
MHEARTBEAT: Pesan heartbeat yang menunjukkan tugas sinkronisasi berjalan normal meskipun tidak ada data baru dari sumber.
timestamp
JSONObject yang berisi timestamp terkait catatan data ini.
-
eventTime: Nilai Long yang merepresentasikan waktu terjadinya perubahan di database sumber. Ini adalah timestamp 13 digit dengan presisi milidetik.
-
systemTime: Nilai Long yang merepresentasikan waktu saat tugas sinkronisasi memproses pesan perubahan ini. Ini adalah timestamp 13 digit dengan presisi milidetik.
-
checkpointTime: Nilai Long yang merepresentasikan waktu yang digunakan untuk mengatur ulang offset sinkronisasi. Nilai ini biasanya sama dengan
eventTime. Ini adalah timestamp 13 digit dengan presisi milidetik.
ddl
Bidang ini diisi hanya untuk operasi DDL yang mengubah struktur tabel. Untuk operasi DML seperti penyisipan, penghapusan, dan modifikasi data, bidang ddl bernilai null.
-
text: String yang berisi teks pernyataan DDL database.
-
ddlMeta: String yang berisi representasi biner objek SQLStatement, diencode dalam Base64. Objek SQLStatement dihasilkan dengan menguraikan pernyataan DDL menggunakan FastSQL.
Jika Anda mengaktifkan dukungan DDL, sistem akan meneruskan objek SQLStatement yang telah diserialisasi. Komponen downstream kemudian dapat mendeserialisasi objek ini untuk merekonstruksi pernyataan DDL untuk sumber data tujuan dan menerapkan perubahan tersebut.
version
N/A
Nomor versi format.
-
-
Serialisasi BLOB
Dalam format JSON ini, setiap pesan berkorespondensi dengan satu JSONObject. Struktur JSONObject ini, yang dapat mencakup objek dan array bersarang, mendefinisikan format pesan.
Untuk melakukan serialisasi pesan, konversi JSONObject menjadi string (misalnya, dengan menggunakan metode
toJSONStringdari fastjson), lalu ubah string tersebut menjadi array byte menggunakan metodeString.getBytes(Charsets.UTF_8).
-
Contoh pesan JSON
-
Insert:
{ "schema": { "dataColumn": [ { "name": "id", "type": "LONG" }, { "name": "name", "type": "STRING" }, { "name": "comment", "type": "STRING" } ], "source": { "dbName": "example_db", "dbType": "MySQL", "tableName": "example_table_pk" }, "primaryKey": [ "id", "name" ] }, "payload": { "op": "INSERT", "after": { "dataColumn": { "name": "joe", "comment": "comment", "id": 1 } }, "sequenceId": "1605339516000000004", "timestamp": { "eventTime": 1605339932000, "systemTime": 1605339932736, "checkpointTime": 1605339932000 } }, "version": "0.0.1" } -
Update before:
{ "schema": { "dataColumn": [ { "name": "id", "type": "LONG" }, { "name": "name", "type": "STRING" }, { "name": "comment", "type": "STRING" } ], "source": { "dbName": "example_db", "dbType": "MySQL", "tableName": "example_table_pk" }, "primaryKey": [ "id", "name" ] }, "payload": { "op": "UPDATE_BEFOR", "before": { "dataColumn": { "name": "joe", "comment": "comment", "id": 1 } }, "sequenceId": "1605339516000000005", "timestamp": { "eventTime": 1605339934000, "systemTime": 1605339934951, "checkpointTime": 1605339934000 } }, "version": "0.0.1" } -
Update after:
{ "schema": { "dataColumn": [ { "name": "id", "type": "LONG" }, { "name": "name", "type": "STRING" }, { "name": "comment", "type": "STRING" } ], "source": { "dbName": "example_db", "dbType": "MySQL", "tableName": "example_table_pk" }, "primaryKey": [ "id", "name" ] }, "payload": { "op": "UPDATE_AFTER", "after": { "dataColumn": { "name": "joe", "comment": "com1", "id": 1 } }, "sequenceId": "1605339516000000005", "timestamp": { "eventTime": 1605339934000, "systemTime": 1605339934951, "checkpointTime": 1605339934000 } }, "version": "0.0.1" } -
Delete:
{ "schema": { "dataColumn": [ { "name": "id", "type": "LONG" }, { "name": "name", "type": "STRING" }, { "name": "comment", "type": "STRING" } ], "source": { "dbName": "example_db", "dbType": "MySQL", "tableName": "example_table_pk" }, "primaryKey": [ "id", "name" ] }, "payload": { "op": "DELETE", "before": { "dataColumn": { "name": "joe", "comment": "com1", "id": 1 } }, "sequenceId": "1605339516000000006", "timestamp": { "eventTime": 1605339937000, "systemTime": 1605339937671, "checkpointTime": 1605339937000 } }, "version": "0.0.1" } -
Heartbeat:
{ "schema": {}, "payload": { "op": "MHEARTBEAT", "timestamp": { "eventTime": 1605339953629, "checkpointTime": 1605339953629 } }, "version": "0.0.1" } -
DDL:
{ "schema": { "source": { "dbName": "example_db", "dbType": "MySQL", "tableName": "example_table_nopk" } }, "payload": { "op": "ALTER", "sequenceId": "1605339516000000035", "ddl": { "text": "alter table example_table_nopk add column holo text", "ddlMeta": "rO0ABXNyACljb20uYWxpYmFiYS5kaS5wbHVnaW4uY2VudGVyLm1ldGEuRERMTWV0YQLb5Cx/YWXtAgACTAAHZGRsVGV4dHQAEkxqYXZhL2xhbmcvU3RyaW5nO0wACXN0YXRlbWVudHQAKkxjb20vYWxpYmFiYS9mYXN0c3FsL3NxbC9hc3QvU1FMU3RhdGVtZW50O3hwdAAtYWx0ZXIgdGFibGUgdF9zaGl5dV9ub3BrIGFkZCBjb2x1bW4gaG9sbyB0ZXh0c3IAPGNvbS5hbGliYWJhLmZhc3RzcWwuc3FsLmFzdC5zdGF0ZW1lbnQuU1FMQWx0ZXJUYWJsZVN0YXRlbWVudBQPP3vMUl2cAgAPSQAHYnVja2V0c1oABmlnbm9yZVoAF2ludmFsaWRhdGVHbG9iYWxJbmRleGVzWgAPbWVyZ2VTbWFsbEZpbGVzWgAHb2ZmbGluZVoABm9ubGluZVoADnJlbW92ZVBhdGl0aW5nWgATdXBkYXRlR2xvYmFsSW5kZXhlc1oAD3VwZ3JhZGVQYXRpdGluZ0wAC2NsdXN0ZXJlZEJ5dAAQTGphdmEvdXRpbC9MaXN0O0wABWl0ZW1zcQB+AAZMAAlwYXJ0aXRpb250ACxMY29tL2FsaWJhYmEvZmFzdHNxbC9zcWwvYXN0L1NRTFBhcnRpdGlvbkJ5O0wACHNvcnRlZEJ5cQB+AAZMAAx0YWJsZU9wdGlvbnNxAH4ABkwAC3RhYmxlU291cmNldAA6TGNvbS9hbGliYWJhL2Zhc3RzcWwvc3FsL2FzdC9zdGF0ZW1lbnQvU1FMRXhwclRhYmxlU291cmNlO3hyACxjb20uYWxpYmFiYS5mYXN0c3FsLnNxbC5hc3QuU1FMU3RhdGVtZW50SW1wbEOxUUDVCJMGAgADWgAJYWZ0ZXJTZW1pTAAGZGJUeXBldAAcTGNvbS9hbGliYWJhL2Zhc3RzcWwvRGJUeXBlO0wACWhlYWRIaW50c3EAfgAGeHIAKWNvbS5hbGliYWJhLmZhc3RzcWwuc3FsLmFzdC5TUUxPYmplY3RJbXBs5LvqLFggFVECAAVJAAxzb3VyY2VDb2x1bW5JAApzb3VyY2VMaW5lTAAKYXR0cmlidXRlc3QAD0xqYXZhL3V0aWwvTWFwO0wABGhpbnR0ACxMY29tL2FsaWJhYmEvZmFzdHNxbC9zcWwvYXN0L1NRTENvbW1lbnRIaW50O0wABnBhcmVudHQAJ0xjb20vYWxpYmFiYS9mYXN0c3FsL3NxbC9hc3QvU1FMT2JqZWN0O3hwAAAAAAAAAABwcHAAfnIAGmNvbS5hbGliYWJhLmZhc3RzcWwuRGJUeXBlAAAAAAAAAAASAAB4cgAOamF2YS5sYW5nLkVudW0AAAAAAAAAABIAAHhwdAAFbXlzcWxwAAAAAAAAAAAAAAAAc3IAE2phdmEudXRpbC5BcnJheUxpc3R4gdIdmcdhnQMAAUkABHNpemV4cAAAAAB3BAAAAAB4c3EAfgAUAAAAAXcEAAAAAXNyADxjb20uYWxpYmFiYS5mYXN0c3FsLnNxbC5hc3Quc3RhdGVtZW50LlNRTEFsdGVyVGFibGVBZGRDb2x1bW4l5T6CFe//BAIABloAB2Nhc2NhZGVaAAVmaXJzdEwAC2FmdGVyQ29sdW1udAAlTGNvbS9hbGliYWJhL2Zhc3RzcWwvc3FsL2FzdC9TUUxOYW1lO0wAB2NvbHVtbnNxAH4ABkwAC2ZpcnN0Q29sdW1ucQB+ABhMAAhyZXN0cmljdHQAE0xqYXZhL2xhbmcvQm9vbGVhbjt4cQB+AAsAAAAAAAAAAHBwcQB+AA8AAHBzcQB+ABQAAAABdwQAAAABc3IAOWNvbS5hbGliYWJhLmZhc3RzcWwuc3FsLmFzdC5zdGF0ZW1lbnQuU1FMQ29sdW1uRGVmaW5pdGlvbst0gLKZ0qAtAgAmWgANYXV0b0luY3JlbWVudFoADGRpc2FibGVJbmRleFoAB3ByZVNvcnRJAAxwcmVTb3J0T3JkZXJaAAZzdG9yZWRaAAd2aXJ0dWFsWgAHdmlzaWJsZUwACGFubkluZGV4dAApTGNvbS9hbGliYWJhL2Zhc3RzcWwvc3FsL2FzdC9TUUxBbm5JbmRleDtMAAZhc0V4cHJ0ACVMY29tL2FsaWJhYmEvZmFzdHNxbC9zcWwvYXN0L1NRTEV4cHI7TAALY2hhcnNldEV4cHJxAH4AHkwADWNvbFByb3BlcnRpZXNxAH4ABkwAC2NvbGxhdGVFeHBycQB+AB5MAAdjb21tZW50cQB+AB5MAAtjb21wcmVzc2lvbnQALkxjb20vYWxpYmFiYS9mYXN0c3FsL3NxbC9hc3QvZXhwci9TUUxDaGFyRXhwcjtMAAtjb25zdHJhaW50c3EAfgAGTAAIZGF0YVR5cGV0AClMY29tL2FsaWJhYmEvZmFzdHNxbC9zcWwvYXN0L1NRTERhdGFUeXBlO0wABmRiVHlwZXEAfgAKTAALZGVmYXVsdEV4cHJxAH4AHkwACWRlbGltaXRlcnEAfgAeTAASZGVsaW1pdGVyVG9rZW5pemVycQB+AB5MAAZlbmFibGVxAH4AGUwABmVuY29kZXEAfgAfTAAGZm9ybWF0cQB+AB5MABBnZW5lcmF0ZWRBbGF3c0FzcQB+AB5MAAhpZGVudGl0eXQARExjb20vYWxpYmFiYS9mYXN0c3FsL3NxbC9hc3Qvc3RhdGVtZW50L1NRTENvbHVtbkRlZmluaXRpb24kSWRlbnRpdHk7TAASanNvbkluZGV4QXR0cnNFeHBycQB+AB5MAAhtYXBwZWRCeXEAfgAGTAAEbmFtZXEAfgAYTAAMbmxwVG9rZW5pemVycQB+AB5MAAhvblVwZGF0ZXEAfgAeTAAEcmVseXEAfgAZTAAMc2VxdWVuY2VUeXBldAAvTGNvbS9hbGliYWJhL2Zhc3RzcWwvc3FsL2FzdC9BdXRvSW5jcmVtZW50VHlwZTtMAARzdGVwcQB+AB5MAAdzdG9yYWdlcQB+AB5MAAl1bml0Q291bnRxAH4AHkwACXVuaXRJbmRleHEAfgAeTAAIdmFsaWRhdGVxAH4AGUwACXZhbHVlVHlwZXEAfgAeeHEAfgALAAAAAAAAAABwcHEAfgAaAAAAAAAAAAAAAHBwcHBwcHBzcQB+ABQAAAAAdwQAAAAAeHNyADpjb20uYWxpYmFiYS5mYXN0c3FsLnNxbC5hc3Quc3RhdGVtZW50LlNRTENoYXJhY3RlckRhdGFUeXBlqtJac/d+04cCAAVaAAloYXNCaW5hcnlMAAtjaGFyU2V0TmFtZXEAfgABTAAIY2hhclR5cGVxAH4AAUwAB2NvbGxhdGVxAH4AAUwABWhpbnRzcQB+AAZ4cgArY29tLmFsaWJhYmEuZmFzdHNxbC5zcWwuYXN0LlNRTERhdGFUeXBlSW1wbEWL29pc1gZFAgAJSgAObmFtZUhhc2hDb2RlNjRaAAh1bnNpZ25lZFoAEXdpdGhMb2NhbFRpbWVab25lWgAIemVyb2ZpbGxMAAlhcmd1bWVudHNxAH4ABkwABmRiVHlwZXEAfgAKTAAHaW5kZXhCeXEAfgAeTAAEbmFtZXEAfgABTAAMd2l0aFRpbWVab25lcQB+ABl4cQB+AAsAAAAAAAAAAHBwcQB+ACP6BPTvGZVAfgAAAHNxAH4AFAAAAAB3BAAAAAB4cHB0AAR0ZXh0cABwcHBwcQB+ABJwcHBwcHBwcHBwc3IAMmNvbS5hbGliYWJhLmZhc3RzcWwuc3FsLmFzdC5leHByLlNRTElkZW50aWZpZXJFeHBy3DXH1zvWbgkCAARKAApoYXNoQ29kZTY0TAAEbmFtZXEAfgABTAAOcmVzb2x2ZWRDb2x1bW5xAH4ADkwAE3Jlc29sdmVkT3duZXJPYmplY3RxAH4ADnhyACdjb20uYWxpYmFiYS5mYXN0c3FsLnNxbC5hc3QuU1FMRXhwckltcGxs2ypmFJxWrQIAAHhxAH4ACwAAAAAAAAAAcHBwQCnxzH5tIDl0AAx0X3NoaXl1X25vcGtwcHBwcA==" }, "timestamp": { "eventTime": 1605342109000, "systemTime": 1605342109259, "checkpointTime": 1605342109000 } }, "version": "0.0.1" }