Este tópico descreve a estrutura e os campos das mensagens gravadas no Kafka.
Contexto
Ao sincronizar um banco de dados completo com o Kafka, a tarefa de sincronização lê dados de uma fonte upstream e os grava em um tópico do Kafka no formato JSON descrito neste documento. O formato geral da mensagem inclui informações das colunas do registro de alteração e o estado dos dados antes e depois da mudança. Para que os consumidores acompanhem o progresso da tarefa, a sincronização também envia periodicamente um registro de heartbeat ao tópico do Kafka. Esse registro possui um campo op com o valor MHEARTBEAT. Este tópico descreve o formato geral da mensagem, o formato da mensagem de heartbeat e os formatos de mensagem para alterações de dados na origem. Para mais informações, consulte Tipos de campo e Parâmetros.
Formato de mensagem para saída de tabela única em tempo real
Ao configurar um destino Kafka para uma tarefa de sincronização de tabela única em tempo real, especifique o formato de valor dos registros gravados no Kafka. Os formatos suportados incluem Canal CDC e JSON. Para mais detalhes, consulte Apêndice: Descrição do formato de saída.
Tipos de campo
O sistema lê os dados da origem, mapeia-os para um dos seis tipos (BOOLEAN, DOUBLE, DATE, BYTES, LONG ou STRING) e os grava em um tópico do Kafka no formato JSON.
|
Tipo |
Descrição |
|
BOOLEAN |
Corresponde ao tipo boolean do JSON. Valores válidos: true e false. |
|
DATE |
Corresponde ao tipo number do JSON. O valor é um timestamp Unix de 13 dígitos com precisão de milissegundos (ms). |
|
BYTES |
Corresponde ao tipo string do JSON. Antes da gravação no Kafka, o array de bytes é codificado em Base64 como string. Durante o consumo, decodifique os dados de Base64 (Codificação: Base64.getEncoder().encodeToString(text.getBytes("UTF-8")); Decodificação: Base64.getDecoder().decode(encodedText)). |
|
STRING |
Corresponde ao tipo string do JSON. |
|
LONG |
Corresponde ao tipo number do JSON. |
|
DOUBLE |
Corresponde ao tipo number do JSON. |
Parâmetros
A seção a seguir descreve os campos presentes nas mensagens gravadas no Kafka.
|
Elemento de nível 1 |
Elemento de nível 2 |
Descrição |
|
schema |
dataColumn |
Tipo JSONArray. Contém as informações de tipo das colunas de dados. O campo dataColumn registra todas as colunas e seus respectivos tipos para os registros de alteração de dados upstream. As alterações incluem modificações de dados (inserções, exclusões e atualizações) e mudanças no schema da tabela.
|
|
primaryKey |
Tipo List. Contém informações da chave primária. pk: Nome da chave primária. |
|
|
source |
Tipo Object. Contém informações do banco de dados ou tabela de origem.
|
|
|
payload |
before |
Tipo JSONObject. Dados antes da alteração. Por exemplo, se a origem for MySQL e ocorrer uma atualização em um registro, o campo before armazena o conteúdo anterior à atualização.
|
|
after |
Dados após a alteração. Segue o mesmo formato do campo before. |
|
|
sequenceId |
Tipo String. Gerado pelo Streamx. Utilizado para ordenar dados durante a mesclagem de dados incrementais e completos. Cada registro do Streamx possui um sequenceId único. Nota
Para uma mensagem de operação de atualização lida da origem, dois registros de gravação são gerados: um registro de update before e um de update after. Ambos compartilham o mesmo sequenceId. |
|
|
scn |
Válido apenas quando a origem é um banco de dados Oracle. Corresponde às informações de SCN do Oracle. |
|
|
op |
Tipo de mensagem lida da origem. Valores válidos:
|
|
|
timestamp |
Tipo JSONObject. Contém timestamps relacionados a este registro.
|
|
|
ddl |
Preenchido apenas quando há alteração no schema da tabela. Para mudanças de dados (incluindo inserções, exclusões e atualizações), o campo ddl é definido como null.
|
|
|
version |
N/A |
Número da versão do formato. |
Formato geral da mensagem
O formato geral das mensagens gravadas no Kafka é o seguinte:
{ "schema": { //Metadados da alteração, contendo apenas nomes e tipos de colunas "dataColumn": [//Informações das colunas dos dados alterados, usadas para atualize registros da tabela de destino { "name": "id", "type": "LONG" }, { "name": "name", "type": "STRING" }, { "name": "binData", "type": "BYTES" }, { "name": "ts", "type": "DATE" }, { "name":"rowid",// Quando a fonte de dados é Oracle, rowid é incluído nas colunas de dados "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"//Tipo String, informações de rowid do Oracle } }, "after": { "dataColumn":{ "id": 222, "name":"donald", "binData": "[base64 string]", "ts": 1590315269000, "rowid": "AAIUMPAAFAACxExAAE"//Tipo String, informações de rowid do Oracle } }, "sequenceId":"XXX",//Tipo String, usado para ordenar dados durante a mesclagem de dados incrementais e completos "scn":"xxxx",//Tipo String, informações de SCN do Oracle "op": "INSERT/UPDATE_BEFOR/UPDATE_AFTER/DELETE/TRANSACTION_BEGIN/TRANSACTION_END/CREATE/ALTER/ERASE/QUERY/TRUNCATE/RENAME/CINDEX/DINDEX/GTID/XACOMMIT/XAROLLBACK/MHEARTBEAT...",//Sensível a maiúsculas e minúsculas "timestamp": { "eventTime": 1,//Obrigatório. O momento em que a alteração ocorreu na origem. Timestamp de 13 dígitos com precisão de milissegundos. "systemTime": 2,//Opcional. O momento em que a tarefa de sincronização processou esta mensagem de alteração. Timestamp de 13 dígitos com precisão de milissegundos. "checkpointTime": 3//Opcional. O momento definido quando o checkpoint de sincronização é redefinido. Timestamp de 13 dígitos com precisão de milissegundos, geralmente igual a eventTime. }, "ddl": { "text": "ADD COLUMN ...", "ddlMeta": "[Binário serializado de SQLStatement, expresso em string base64]" } }, "version":"1.0.0" }
Para mais informações sobre tipos de campo e parâmetros, consulte Tipos de campo e Parâmetros.
Formato da mensagem 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"
}
Para mais informações sobre tipos de campo e parâmetros, consulte Tipos de campo e Parâmetros.
Formatos de mensagem para alterações de dados na origem
-
Formato de mensagem do Kafka para inserções de dados na origem:
{ "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" } -
Formato de mensagem do Kafka para atualizações de dados na origem:
-
Se a opção When one record in the source is updated, one Kafka record is generated. não estiver selecionada, o formato de mensagem do Kafka para uma atualização de dados na origem consistirá em duas mensagens que descrevem, respectivamente, o estado dos dados antes e depois da atualização. Os formatos são:
Formato de mensagem para o estado dos dados antes da atualização:
{ "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" }Formato de mensagem para o estado dos dados após a atualização:
{ "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" } -
Caso a opção When one record in the source is updated, one Kafka record is generated. esteja selecionada, o formato de mensagem do Kafka para uma atualização de dados na origem consistirá em uma única mensagem que descreve tanto o estado dos dados antes quanto depois da atualização. O formato é:
{ "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" }
-
-
Formato de mensagem do Kafka para exclusões de dados na origem:
{ "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" }
Para mais informações sobre tipos de campo e parâmetros, consulte Tipos de campo e Parâmetros.