すべてのプロダクト
Search
ドキュメントセンター

DataWorks:付録:メッセージフォーマット

最終更新日:May 05, 2026

このトピックでは、Kafka に書き込まれるメッセージの構造とフィールドについて説明します。

背景情報

データベース全体を Kafka に同期する場合、同期タスクは上流のデータソースからデータを読み取り、このトピックで説明する JSON フォーマットで Kafka トピックに書き込みます。メッセージの全体的なフォーマットには、変更記録のカラム情報と、変更前後のデータ状態が含まれます。コンシューマーがタスクの進捗を追跡できるように、同期タスクは定期的にハートビートレコードを Kafka トピックに送信します。このハートビートレコードには、値が MHEARTBEATop フィールドが含まれます。このトピックでは、メッセージの全体フォーマットハートビートメッセージフォーマット、およびソースデータ変更時のメッセージフォーマットについて説明します。詳細については、「フィールドの型」および「パラメーター」をご参照ください。

リアルタイム単一テーブル出力のメッセージフォーマット

リアルタイム単一テーブル同期タスクの Kafka 送信先を設定する際は、Kafka に書き込まれるレコードの値フォーマットを指定する必要があります。サポートされるフォーマットには、Canal CDC と JSON があります。詳細については、「付録:出力フォーマットの説明」をご参照ください。

フィールドの型

システムはソースからデータを読み取り、次の 6 種類のいずれか(BOOLEAN、DOUBLE、DATE、BYTES、LONG、STRING)にマッピングして、JSON フォーマットで Kafka トピックに書き込みます。

説明

BOOLEAN

JSON のブール型に対応します。有効な値は true および false です。

DATE

JSON の数値型に対応します。値は、ミリ秒 (ms) 精度の 13 桁の UNIX タイムスタンプです。

BYTES

JSON の文字列型に対応します。Kafka に書き込まれる前に、バイト配列は Base64 エンコードされて文字列になります。消費時には、Base64 デコードする必要があります(エンコード:Base64.getEncoder().encodeToString(text.getBytes("UTF-8"))。デコード:Base64.getDecoder().decode(encodedText))。

STRING

JSON の文字列型に対応します。

LONG

JSON の数値型に対応します。

DOUBLE

JSON の数値型に対応します。

パラメーター

以下の表は、Kafka に書き込まれるメッセージ内の各フィールドについて説明しています。

レベル 1 フィールド

レベル 2 フィールド

説明

schema

dataColumn

データ型は JSONArray です。dataColumn パラメーターは、上流のデータ変更レコードのすべてのカラムに関する型情報を記録します。変更操作には、データベース内のデータの追加・削除・更新などのデータ変更や、データベーステーブルスキーマの変更が含まれます。

  • name:カラム名。

  • type:カラムの型。

primaryKey

プライマリキー列名のリスト。

source

ソースデータベースまたはテーブルに関する情報を含むオブジェクト。

  • dbType:データベースの型。文字列。

  • dbVersion:データベースのバージョン。文字列。

  • dbName:データベース名。文字列。

  • schemaName:スキーマ名。PostgreSQL や SQL Server などのデータベースで使用されます。文字列。

  • tableName:テーブル名。文字列。

payload

before

変更前のデータレコードを表す JSONObject。別名「before イメージ」とも呼ばれます。たとえば、MySQL データソースのレコードが更新された場合、before フィールドには更新前のレコード内容が格納されます。

  • このフィールドは、更新および削除操作に対して入力されます。

  • dataColumn:データ情報。型は JSONObject です。形式は カラム名:カラム値 で、カラム名は文字列、カラム値は BOOLEAN、DOUBLE、DATE、BYTES、LONG、STRING のいずれかです。

after

変更後のデータ。別名「after イメージ」とも呼ばれます。形式は before フィールドと同じです。

sequenceId

StreamX が各レコードに対して生成する一意の文字列。完全データと増分データをマージする際にデータをソートするために使用されます。

説明

ソースからの更新操作により、2 つのレコード(「update before」レコードと「update after」レコード)が生成されます。両方のレコードは同じ sequenceId を共有します。

scn

ソースからのシステム変更番号(SCN)。このフィールドは Oracle ソースにのみ適用されます。

op

ソースからキャプチャされた操作の種類。有効な値:

  • INSERT:データの挿入。

  • UPDATE_BEFOR:更新の before イメージ。

  • UPDATE_AFTER:更新の after イメージ。

  • DELETE:データの削除。

  • TRANSACTION_BEGIN:データベーストランザクションの開始。

  • TRANSACTION_END:データベーストランザクションの終了。

  • CREATE:テーブルの作成。

  • ALTER:テーブルの変更。

  • QUERY:データベース変更の元となる SQL ステートメント。

  • TRUNCATE:テーブルの切り捨て。

  • RENAME:テーブルの名前変更操作。

  • CINDEX:インデックスの作成。

  • DINDEX:インデックスの削除。

  • MHEARTBEAT:ソースに新しいデータがなくても同期タスクがアクティブであることを示すハートビートメッセージ。

timestamp

データレコードに関連するタイムスタンプを含む JSONObject。

  • eventTime:ソースデータベースで変更が発生した時刻。ミリ秒精度の 13 桁のタイムスタンプ。Long 型。

  • systemTime:同期タスクが変更メッセージを処理した時刻。ミリ秒精度の 13 桁のタイムスタンプ。Long 型。

  • checkpointTime:同期オフセットをリセットする際に使用されるタイムスタンプ。通常は eventTime と同じ値です。ミリ秒精度の 13 桁のタイムスタンプ。Long 型。

ddl

このフィールドは、データベースのテーブルスキーマが変更された場合にのみデータが入力されます。データの追加・削除・変更などのデータ変更では、対応する ddl は null に設定されます。

  • text:DDL ステートメントのテキスト。文字列。

  • ddlMeta:DDL 変更を記録したシリアル化された Java オブジェクトの Base64 エンコード文字列。文字列。

version

N/A

メッセージフォーマットのバージョン。

メッセージの全体フォーマット

以下のコードは、Kafka メッセージの全体フォーマットを示しています。

{
    "schema": { // 変更に関するメタデータ(カラム名と型など)。
        "dataColumn": [// 変更されたデータカラムの情報。送信先のレコード更新に使用。
            {
                "name": "id",
                "type": "LONG"
            },
            {
                "name": "name",
                "type": "STRING"
            },
            {
                "name": "binData",
                "type": "BYTES"
            },
            {
                "name": "ts",
                "type": "DATE"
            },
            {
              "name":"rowid",// データソースが Oracle の場合、rowid がデータカラムに含まれます。
              "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"// 文字列型。Oracle ソースからの rowid。
            }
        },
        "after": {
            "dataColumn":{
                "id": 222,
                "name":"donald",
                "binData": "[base64 string]",
                "ts": 1590315269000,
                "rowid": "AAIUMPAAFAACxExAAE"// 文字列型。Oracle ソースからの rowid。
            }
        },
        "sequenceId":"XXX",// 文字列型。完全データと増分データをマージする際にデータをソートするために使用。
        "scn":"xxxx",// 文字列型。Oracle ソースからのシステム変更番号(SCN)。
        "op": "INSERT/UPDATE_BEFOR/UPDATE_AFTER/DELETE/TRANSACTION_BEGIN/TRANSACTION_END/CREATE/ALTER/ERASE/QUERY/TRUNCATE/RENAME/CINDEX/DINDEX/GTID/XACOMMIT/XAROLLBACK/MHEARTBEAT...",// 値は大文字小文字を区別します。
        "timestamp": {
            "eventTime": 1,// 必須。ソースデータベースでの変更イベントの時刻。ミリ秒精度の 13 桁のタイムスタンプ。
            "systemTime": 2,// 任意。同期タスクがこの変更メッセージを処理した時刻。ミリ秒精度の 13 桁のタイムスタンプ。
            "checkpointTime": 3// 任意。同期オフセットをリセットする際に使用されるタイムスタンプ。通常は eventTime と同じ。
        },
        "ddl": {
            "text": "ADD COLUMN ...",
            "ddlMeta": "[シリアライズされた SQLStatement オブジェクトの Base64 エンコード文字列]"
        }
    },
    "version":"1.0.0"
}
説明

フィールドの型およびパラメーターの詳細については、「フィールドの型」および「パラメーター」をご参照ください。

ハートビートメッセージフォーマット

{
    "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"
}
説明

フィールドの型およびパラメーターの詳細については、「フィールドの型」および「パラメーター」をご参照ください。

ソースデータ変更時のメッセージフォーマット

  • 挿入操作のメッセージフォーマット:

    {
        "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"
    }
  • 更新操作のメッセージフォーマット:

    • When one record in the source is updated, one Kafka record is generated.」オプションが選択されていない場合、更新操作により 2 つの Kafka メッセージが生成されます。1 つは更新前のデータ状態(「before イメージ」)、もう 1 つは更新後のデータ状態(「after イメージ」)です。以下のサンプルコードはそのフォーマットを示しています。

      before イメージのメッセージフォーマット:

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

      after イメージのメッセージフォーマット:

      {
          "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"
      }
    • When one record in the source is updated, one Kafka record is generated.」オプションが選択されている場合、更新操作により before イメージと after イメージの両方を含む 1 つの Kafka メッセージが生成されます。以下のサンプルコードはそのフォーマットを示しています。

      {
          "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"
      }
  • 削除操作のメッセージフォーマット:

    {
        "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"
    }
説明

フィールドの型およびパラメーターの詳細については、「フィールドの型」および「パラメーター」をご参照ください。