Tous les produits
Search
Centre de documentation

DataWorks:Annexe : Format des messages

Dernière mise à jour :Aug 10, 2026

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 true et false.

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.

  • name : nom de la colonne.

  • type : type de la colonne.

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.

  • dbType : type String. Type de base de données.

  • dbVersion : type String. Version de la base de données.

  • dbName : type String. Nom de la base de données.

  • schemaName : type String. Nom du schéma (pour les bases de données telles que PostgreSQL et SQL Server).

  • tableName : type String. Nom de la table.

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 before stocke le contenu des données avant la mise à jour.

  • Ce champ est renseigné lorsqu’un message de mise à jour ou de suppression est lu depuis la source.

  • dataColumn : type JSONObject. Représente les informations de données. Le format est nom de la colonne : valeur de la colonne. Le nom de la colonne est une chaîne et la valeur de la colonne est de type BOOLEAN, DOUBLE, DATE, BYTES, LONG ou STRING.

after

Données après la modification. Le format est identique à celui de before.

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 sequenceId unique.

Remarque

Pour un message de mise à jour lu depuis la source, deux enregistrements d’écriture sont générés : un enregistrement before de mise à jour et un enregistrement after de mise à jour. Ces deux enregistrements partagent le même sequenceId.

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 :

  • INSERT : insertion de données.

  • UPDATE_BEFOR : données avant la mise à jour.

  • UPDATE_AFTER : données après la mise à jour.

  • DELETE : suppression de données.

  • TRANSACTION_BEGIN : début de transaction de base de données.

  • TRANSACTION_END : fin de transaction de base de données.

  • CREATE : création de table de base de données.

  • ALTER : modification de table de base de données.

  • QUERY : SQL brut de la modification de base de données.

  • TRUNCATE : troncature de table de base de données.

  • RENAME : renommage de table de base de données.

  • CINDEX : création d’index.

  • DINDEX : suppression d’index.

  • MHEARTBEAT : message de heartbeat indiquant que la synchronisation est toujours active en l’absence de nouvelles données générées à la source.

timestamp

Type JSONObject. Contient les horodatages associés à cet enregistrement.

  • eventTime : type Long. Heure à laquelle la modification s’est produite à la source. La valeur est un horodatage à 13 chiffres avec une précision en millisecondes.

  • systemTime : type Long. Heure à laquelle la tâche de synchronisation a traité ce message de modification. La valeur est un horodatage à 13 chiffres avec une précision en millisecondes.

  • checkpointTime : type Long. Heure définie lors de la réinitialisation du point de contrôle de synchronisation. La valeur est un horodatage à 13 chiffres avec une précision en millisecondes et est généralement identique à eventTime.

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 null.

  • text : type String. Texte de l’instruction DDL.

  • ddlMeta : type String. Chaîne encodée en Base64 obtenue par sérialisation d’un objet Java enregistrant la modification DDL.

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

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

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

Pour plus d’informations sur les types de champs et les paramètres, consultez les sections Types de champs et Paramètres.