Cette rubrique décrit la structure et les champs des messages écrits dans Kafka.
Contexte
Lors de la synchronisation d’une base de données complète vers Kafka, la tâche lit les données depuis une source en amont et les écrit dans un topic Kafka au format JSON décrit dans cette rubrique. Le format global du message inclut les informations de colonne pour l’enregistrement de modification, ainsi que l’état des données avant et après la modification. Pour permettre aux consommateurs de suivre la progression de la tâche, celle-ci envoie également périodiquement un enregistrement de type heartbeat au topic Kafka. Cet enregistrement contient un champ op dont la valeur est MHEARTBEAT. Cette rubrique décrit le format global des messages, le format des messages de heartbeat et les formats des messages liés aux modifications des données sources. Pour plus d’informations, consultez les sections Types de champs et Paramètres.
Format des messages pour la sortie en temps réel d’une table unique
Lorsque vous configurez une destination Kafka pour une tâche de synchronisation en temps réel d’une table unique, vous devez spécifier le format des valeurs des enregistrements écrits dans Kafka. Les formats pris en charge sont Canal CDC et JSON. Pour plus de détails, consultez l’Annexe : Description du format de sortie.
Types de champs
Le système lit les données depuis la source, les mappe vers l’un des six types (BOOLEAN, DOUBLE, DATE, BYTES, LONG ou STRING), puis les écrit dans un topic Kafka au format JSON.
|
Type |
Description |
|
BOOLEAN |
Correspond au type booléen JSON. Les valeurs valides sont |
|
DATE |
Correspond au type numérique JSON. La valeur est un horodatage Unix à 13 chiffres, avec une précision en millisecondes (ms). |
|
BYTES |
Correspond au type chaîne JSON. Avant l’écriture dans Kafka, le tableau d’octets est encodé en Base64 sous forme de chaîne. Lors de la consommation, les données doivent être décodées en Base64 (Encodage : Base64.getEncoder().encodeToString(text.getBytes("UTF-8")) ; Décodage : Base64.getDecoder().decode(encodedText)). |
|
STRING |
Correspond au type chaîne JSON. |
|
LONG |
Correspond au type numérique JSON. |
|
DOUBLE |
Correspond au type numérique JSON. |
Paramètres
La section suivante décrit les champs présents dans les messages écrits dans Kafka.
|
Élément de niveau 1 |
Élément de niveau 2 |
Description |
|
schema |
dataColumn |
Type JSONArray. Contient les informations de type pour les colonnes de données. dataColumn répertorie toutes les colonnes et leurs types correspondants pour les enregistrements de modification des données en amont. Les modifications incluent les opérations sur les données (insertions, suppressions et mises à jour) ainsi que les changements de schéma de table.
|
|
primaryKey |
Type List. Contient les informations de clé primaire. pk : nom de la clé primaire. |
|
|
source |
Type Object. Contient les informations sur la base de données ou la table source.
|
|
|
payload |
before |
Type JSONObject. Données avant la modification. Par exemple, si la source est MySQL et qu’une opération de mise à jour est effectuée sur un enregistrement, le champ
|
|
after |
Données après la modification. Le format est identique à celui de |
|
|
sequenceId |
Type String. Généré par Streamx. Utilisé pour trier les données lors de la fusion des données incrémentielles et complètes. Chaque enregistrement Streamx possède un Remarque
Pour un message de mise à jour lu depuis la source, deux enregistrements d’écriture sont générés : un enregistrement |
|
|
scn |
Valide uniquement lorsque la source est une base de données Oracle. Correspond aux informations SCN d’Oracle. |
|
|
op |
Type de message lu depuis la source. Valeurs valides :
|
|
|
timestamp |
Type JSONObject. Contient les horodatages associés à cet enregistrement.
|
|
|
ddl |
Ce champ est renseigné uniquement lorsque le schéma de la table est modifié. Pour les modifications de données (y compris les insertions, suppressions et mises à jour), le champ ddl est défini sur
|
|
|
version |
N/A |
Numéro de version du format. |
Format global des messages
Le format global des messages écrits dans Kafka est le suivant :
{
"schema": { //Metadata of the change, containing only column names and column types
"dataColumn": [//Column information for the changed data, used to update target table records
{
"name": "id",
"type": "LONG"
},
{
"name": "name",
"type": "STRING"
},
{
"name": "binData",
"type": "BYTES"
},
{
"name": "ts",
"type": "DATE"
},
{
"name":"rowid",// When the data source is Oracle, rowid is included in the data columns
"type":"STRING"
}
],
"primaryKey": [
"pkName1",
"pkName2"
],
"source": {
"dbType": "mysql",
"dbVersion": "1.0.0",
"dbName": "myDatabase",
"schemaName": "mySchema",
"tableName": "tableName"
}
},
"payload": {
"before": {
"dataColumn":{
"id": 111,
"name":"scooter",
"binData": "[base64 string]",
"ts": 1590315269000,
"rowid": "AAIUMPAAFAACxExAAE"//String type, Oracle rowid information
}
},
"after": {
"dataColumn":{
"id": 222,
"name":"donald",
"binData": "[base64 string]",
"ts": 1590315269000,
"rowid": "AAIUMPAAFAACxExAAE"//String type, Oracle rowid information
}
},
"sequenceId":"XXX",//String type, used for sorting data during the merge of incremental and full data
"scn":"xxxx",//String type, Oracle SCN information
"op": "INSERT/UPDATE_BEFOR/UPDATE_AFTER/DELETE/TRANSACTION_BEGIN/TRANSACTION_END/CREATE/ALTER/ERASE/QUERY/TRUNCATE/RENAME/CINDEX/DINDEX/GTID/XACOMMIT/XAROLLBACK/MHEARTBEAT...",//Case-sensitive
"timestamp": {
"eventTime": 1,//Required. The time when the change occurred at the source. A 13-digit timestamp with millisecond precision.
"systemTime": 2,//Optional. The time when the synchronization task processed this change message. A 13-digit timestamp with millisecond precision.
"checkpointTime": 3//Optional. The time set when the synchronization checkpoint is reset. A 13-digit timestamp with millisecond precision, generally equal to eventTime.
},
"ddl": {
"text": "ADD COLUMN ...",
"ddlMeta": "[SQLStatement serialized binary, expressed in base64 string]"
}
},
"version":"1.0.0"
}
Pour plus d’informations sur les types de champs et les paramètres, consultez les sections Types de champs et Paramètres.
Format des messages de heartbeat
{
"schema": {
"dataColumn": null,
"primaryKey": null,
"source": null
},
"payload": {
"before": null,
"after": null,
"sequenceId": null,
"timestamp": {
"eventTime": 1620457659000,
"checkpointTime": 1620457659000
},
"op": "MHEARTBEAT",
"ddl": null
},
"version": "0.0.1"
}
Pour plus d’informations sur les types de champs et les paramètres, consultez les sections Types de champs et Paramètres.
Formats des messages liés aux modifications des données sources
-
Format des messages Kafka pour les insertions de données à la source :
{ "schema": { "dataColumn": [ { "name": "name", "type": "STRING" }, { "name": "job", "type": "STRING" }, { "name": "sex", "type": "STRING" }, { "name": "#alibaba_rds_row_id#", "type": "LONG" } ], "primaryKey": null, "source": { "dbType": "MySQL", "dbName": "pkset_test", "tableName": "pkset_test_no_pk" } }, "payload": { "before": null, "after": { "dataColumn": { "name": "name11", "job": "job11", "sex": "man", "#alibaba_rds_row_id#": 15 } }, "sequenceId": "1620457642589000000", "timestamp": { "eventTime": 1620457896000, "systemTime": 1620457896977, "checkpointTime": 1620457896000 }, "op": "INSERT", "ddl": null }, "version": "0.0.1" } -
Format des messages Kafka pour les mises à jour de données à la source :
-
Si l’option When one record in the source is updated, one Kafka record is generated. n’est pas sélectionnée, le format des messages Kafka pour une mise à jour de données source se compose de deux messages Kafka décrivant respectivement l’état des données avant et après la mise à jour. Les formats des messages sont les suivants :
Format du message pour l’état des données avant la mise à jour :
{ "schema": { "dataColumn": [ { "name": "name", "type": "STRING" }, { "name": "job", "type": "STRING" }, { "name": "sex", "type": "STRING" }, { "name": "#alibaba_rds_row_id#", "type": "LONG" } ], "primaryKey": null, "source": { "dbType": "MySQL", "dbName": "pkset_test", "tableName": "pkset_test_no_pk" } }, "payload": { "before": { "dataColumn": { "name": "name11", "job": "job11", "sex": "man", "#alibaba_rds_row_id#": 15 } }, "after": null, "sequenceId": "1620457642589000001", "timestamp": { "eventTime": 1620458077000, "systemTime": 1620458077779, "checkpointTime": 1620458077000 }, "op": "UPDATE_BEFOR", "ddl": null }, "version": "0.0.1" }Format du message pour l’état des données après la mise à jour :
{ "schema": { "dataColumn": [ { "name": "name", "type": "STRING" }, { "name": "job", "type": "STRING" }, { "name": "sex", "type": "STRING" }, { "name": "#alibaba_rds_row_id#", "type": "LONG" } ], "primaryKey": null, "source": { "dbType": "MySQL", "dbName": "pkset_test", "tableName": "pkset_test_no_pk" } }, "payload": { "before": null, "after": { "dataColumn": { "name": "name11", "job": "job11", "sex": "woman", "#alibaba_rds_row_id#": 15 } }, "sequenceId": "1620457642589000001", "timestamp": { "eventTime": 1620458077000, "systemTime": 1620458077779, "checkpointTime": 1620458077000 }, "op": "UPDATE_AFTER", "ddl": null }, "version": "0.0.1" } -
Si l’option When one record in the source is updated, one Kafka record is generated. est sélectionnée, le format des messages Kafka pour une mise à jour de données source se compose d’un seul message Kafka décrivant à la fois l’état des données avant et après la mise à jour. Le format du message est le suivant :
{ "schema": { "dataColumn": [ { "name": "name", "type": "STRING" }, { "name": "job", "type": "STRING" }, { "name": "sex", "type": "STRING" }, { "name": "#alibaba_rds_row_id#", "type": "LONG" } ], "primaryKey": null, "source": { "dbType": "MySQL", "dbName": "pkset_test", "tableName": "pkset_test_no_pk" } }, "payload": { "before": { "dataColumn": { "name": "name11", "job": "job11", "sex": "man", "#alibaba_rds_row_id#": 15 } }, "after": { "dataColumn": { "name": "name11", "job": "job11", "sex": "woman", "#alibaba_rds_row_id#": 15 } }, "sequenceId": "1620457642589000001", "timestamp": { "eventTime": 1620458077000, "systemTime": 1620458077779, "checkpointTime": 1620458077000 }, "op": "UPDATE_AFTER", "ddl": null }, "version": "0.0.1" }
-
-
Format des messages Kafka pour les suppressions de données à la source :
{ "schema": { "dataColumn": [ { "name": "name", "type": "STRING" }, { "name": "job", "type": "STRING" }, { "name": "sex", "type": "STRING" }, { "name": "#alibaba_rds_row_id#", "type": "LONG" } ], "primaryKey": null, "source": { "dbType": "MySQL", "dbName": "pkset_test", "tableName": "pkset_test_no_pk" } }, "payload": { "before": { "dataColumn": { "name": "name11", "job": "job11", "sex": "woman", "#alibaba_rds_row_id#": 15 } }, "after": null, "sequenceId": "1620457642589000002", "timestamp": { "eventTime": 1620458266000, "systemTime": 1620458266101, "checkpointTime": 1620458266000 }, "op": "DELETE", "ddl": null }, "version": "0.0.1" }
Pour plus d’informations sur les types de champs et les paramètres, consultez les sections Types de champs et Paramètres.