Todos os produtos
Search
Central de documentação

DataWorks:Apêndice: Formato de mensagem

Última atualização: Jun 27, 2026

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.

  • name: Nome da coluna.

  • type: Tipo da coluna.

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.

  • dbType: Tipo String. Tipo do banco de dados.

  • dbVersion: Tipo String. Versão do banco de dados.

  • dbName: Tipo String. Nome do banco de dados.

  • schemaName: Tipo String. Nome do schema (para bancos de dados como PostgreSQL e SQL Server).

  • tableName: Tipo String. Nome da tabela.

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.

  • Preenchido quando uma mensagem de operação de atualização ou exclusão é lida da origem.

  • dataColumn: Tipo JSONObject. Informações dos dados. Formato: nome da coluna: valor da coluna. O nome da coluna é uma string e o valor corresponde a um dos tipos BOOLEAN, DOUBLE, DATE, BYTES, LONG ou STRING.

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:

  • INSERT: Inserção de dados.

  • UPDATE_BEFOR: Dados antes da atualização.

  • UPDATE_AFTER: Dados após a atualização.

  • DELETE: Exclusão de dados.

  • TRANSACTION_BEGIN: Início de transação do banco de dados.

  • TRANSACTION_END: Fim de transação do banco de dados.

  • CREATE: Criação de tabela no banco de dados.

  • ALTER: Alteração de tabela no banco de dados.

  • QUERY: SQL bruto da alteração no banco de dados.

  • TRUNCATE: Truncamento de tabela no banco de dados.

  • RENAME: Renomeação de tabela no banco de dados.

  • CINDEX: Criação de índice.

  • DINDEX: Exclusão de índice.

  • MHEARTBEAT: Mensagem de heartbeat indicando que a sincronização continua em execução quando nenhum dado novo é gerado na origem.

timestamp

Tipo JSONObject. Contém timestamps relacionados a este registro.

  • eventTime: Tipo Long. Momento em que a alteração ocorreu na origem. Valor: timestamp de 13 dígitos com precisão de milissegundos.

  • systemTime: Tipo Long. Momento em que a tarefa de sincronização processou esta mensagem de alteração. Valor: timestamp de 13 dígitos com precisão de milissegundos.

  • checkpointTime: Tipo Long. Momento definido quando o checkpoint de sincronização é redefinido. Valor: timestamp de 13 dígitos com precisão de milissegundos, geralmente igual a eventTime.

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.

  • text: Tipo String. Texto da instrução DDL.

  • ddlMeta: Tipo String. String codificada em Base64 obtida pela serialização de um objeto Java que registra a alteração DDL.

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" }

Nota

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"
}
Nota

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"
    }
Nota

Para mais informações sobre tipos de campo e parâmetros, consulte Tipos de campo e Parâmetros.