本記事では、データサブスクリプションでサポートされているデータ型である Avro と Canal JSON のメッセージ形式について説明します。
サブスクライブするメッセージ形式の説明
Avro JSONメッセージ形式
{
"namespace": "com.xxx.oms.xxx.xxx.avro",
"type": "record",
"name": "AvroRecord",
"fields": [
{
"name": "id",
"type": "long",
"doc": "unique id of this record in the whole stream"
},
{
"name": "version",
"default": 1,
"type": [
"int"
],
"doc": "protocol version"
},
{
"name": "operation",
"default": null,
"type": [
"null",
{
"namespace": "com.xxx.oms.xxx.xxx.avro",
"name": "Operation",
"type": "enum",
"symbols": [
"INSERT",
"UPDATE",
"DELETE",
"DDL",
"BEGIN",
"COMMIT",
"HEARTBEAT"
]
}
]
},
{
"name": "xid",
"default": null,
"type": [
"null",
"string"
],
"doc": "transaction id"
},
{
"name": "txind",
"default": null,
"type": [
"null",
{
"namespace": "com.xxx.oms.xxx.xxx.avro",
"name": "TxindType",
"type": "enum",
"symbols": [
"B",
"M",
"E",
"W"
]
}
],
"doc": "transaction record flag"
},
{
"name": "position",
"default": null,
"type": [
"null",
"string"
],
"doc": "current record position"
},
{
"name": "timestamp",
"default": null,
"type": [
"null",
"long"
],
"doc": "record log timestamp"
},
{
"name": "source",
"default": null,
"type": [
"null",
{
"type": "record",
"name": "Source",
"fields": [
{
"name": "sourceType",
"type": {
"type": "enum",
"name": "SourceType",
"symbols": [
"OB_MYSQL",
"OB_ORACLE"
]
}
},
{
"name": "version",
"type": {
"type": "string"
},
"doc": "source datasource version information"
}
]
}
],
"doc": "source datasource"
},
{
"name": "schemaName",
"default": null,
"type": [
"null",
"string"
],
"doc": "schema name"
},
{
"name": "tableName",
"default": null,
"type": [
"null",
"string"
],
"doc": "table name"
},
{
"name": "fields",
"default": null,
"type": [
"null",
{
"type": "array",
"items": {
"namespace": "com.xxx.oms.xxx.xxx.avro",
"name": "Field",
"type": "record",
"fields": [
{
"name": "name",
"type": "string",
"doc": "column name for table"
},
{
"name": "dataTypeNumber",
"type": "int",
"doc": "data type number of this column in source table"
}
],
"doc": "definition of column"
}
}
]
},
{
"name": "pkIndexes",
"default": null,
"type": [
"null",
{
"type": "array",
"items": "int"
}
],
"doc": "primary keys"
},
{
"name": "ukIndexes",
"default": null,
"type": [
"null",
{
"type": "array",
"items": {
"type": "array",
"items": "int"
}
}
],
"doc": "primary keys"
},
{
"name": "beforeImages",
"default": null,
"doc": "values of a row before modify",
"type": [
"null",
{
"type": "array",
"items": {
"type": "record",
"name": "ColumnValue",
"namespace": "com.xxx.oms.xxx.xxx.avro",
"fields": [
{
"name": "type_info",
"type": {
"type": "enum",
"name": "DataType",
"symbols": [
"NULL",
"INTEGER",
"LONG",
"FLOAT",
"DOUBLE",
"DECIMAL",
"STRING",
"BINARY",
"TIMESTAMP",
"DATE",
"TIME",
"DATETIME",
"BOOLEAN",
"BIT",
"GEOMETRY",
"RAW",
"ENUM",
"SET",
"ARRAY",
"VECTOR",
"SPARSE_VECTOR",
"ROARINGBITMAP",
"MAP"
]
}
},
{
"name": "value",
"type": [
"null",
"boolean",
"int",
"long",
"float",
"double",
"bytes",
"string",
{
"type": "record",
"name": "StringObject",
"fields": [
{
"name": "charsetName",
"type": "string"
},
{
"name": "value",
"type": "string"
}
]
},
{
"type": "record",
"name": "DecimalObject",
"fields": [
{
"name": "precision",
"type": "int"
},
{
"name": "scale",
"type": "int"
},
{
"name": "value",
"type": "string"
}
]
},
{
"type": "record",
"name": "DateObject",
"fields": [
{
"name": "year",
"type": "int"
},
{
"name": "month",
"type": "int"
},
{
"name": "day",
"type": "int"
}
]
},
{
"type": "record",
"name": "TimeObject",
"fields": [
{
"name": "negative",
"type": "boolean",
"default": false
},
{
"name": "hours",
"type": "int"
},
{
"name": "minutes",
"type": "int"
},
{
"name": "seconds",
"type": "int"
},
{
"name": "nanos",
"type": "int"
}
]
},
{
"type": "record",
"name": "DateTimeObject",
"fields": [
{
"name": "year",
"type": "int"
},
{
"name": "month",
"type": "int"
},
{
"name": "day",
"type": "int"
},
{
"name": "hours",
"type": "int"
},
{
"name": "minutes",
"type": "int"
},
{
"name": "seconds",
"type": "int"
},
{
"name": "nanos",
"type": "int",
"default": null
}
]
},
{
"type": "record",
"name": "TimestampObject",
"fields": [
{
"name": "seconds",
"type": "long"
},
{
"name": "nanos",
"type": "int"
},
{
"name": "timezone",
"type": [
"null",
"string"
],
"default": null
}
]
},
{
"type": "record",
"name": "BitObject",
"fields": [
{
"name": "bit_length",
"type": "int"
},
{
"name": "value",
"type": "string"
}
]
},
{
"type": "record",
"name": "EnumSetValue",
"fields": [
{
"name": "value",
"type": "string"
},
{
"name": "defines",
"type": [
"null",
{
"type": "array",
"items": "string"
}
]
}
]
},
{
"type": "record",
"name": "GeometryValue",
"fields": [
{
"name": "srid",
"type": "int"
},
{
"name": "wkb",
"type": "bytes"
}
]
}
]
}
]
}
}
]
},
{
"name": "afterImages",
"default": null,
"doc": "values of a row after modify",
"type": [
"null",
{
"type": "array",
"items": "com.xxx.oms.xxx.xxx.avro.ColumnValue"
}
]
},
{
"name": "sql",
"default": null,
"type": [
"null",
"string"
],
"doc": "DDL sql if this is a ddl event"
},
{
"name": "tags",
"default": null,
"type": [
"null",
{
"type": "map",
"values": "string"
}
],
"doc": "Additional information about the source"
},
{
"name": "total",
"type": "int",
"default": -1,
"doc": "Not implemented at current version"
},
{
"name": "index",
"default": -1,
"type": "int",
"doc": "Not implemented at current version"
},
{
"name": "beforeImageBytes",
"default": "",
"type": "bytes",
"doc": "bytes value of beforeImage when record is too large to split"
},
{
"name": "afterImageBytes",
"default": "",
"type": "bytes",
"doc": "bytes value of afterImage when record is too large to split"
}
]
}
パラメータ |
説明 |
|---|---|
| id | グローバルな連番ID。 |
| version | プロトコルバージョン。現在のバージョンは1です。 |
| operation | メッセージタイプ。"INSERT"、"UPDATE"、"DELETE"、"DDL"、"BEGIN"、"COMMIT"、"HEARTBEAT"が指定可能です。 |
| xid | トランザクションID。 |
| txind | トランザクション内での現在のDMLレコードの識別子。
|
| position | 増分ログ内の現在のレコードのオフセット量。 |
| timestamp | 増分レコードの書き込み時間。Unixタイムスタンプ、単位は秒。 |
| source | ソースデータベースのデータベースタイプとバージョン。現在のデータベースタイプはOB_MYSQLとOB_ORACLEのみサポートしています。 |
| schemaName | データベース名またはスキーマ名。 |
| tableName | ソースデータベースのserverId。ソースデータベースのserver_idを確認するには、SHOW VARIABLES LIKE 'server_id'をご参照ください。 |
| fields | テーブルの各列の定義。列名と列タイプを含みます。 |
| pkIndexes | データベースのテーブルに主キーが設定されている場合、DMLメッセージにこのパラメータが含まれます。設定されていない場合は含まれません。 |
| ukIndexes | データベースのテーブルに一意キーが設定されている場合、DMLメッセージにこのパラメータが含まれます。設定されていない場合は含まれません。 |
| beforeImages | DML実行前のその行のデータ。INSERTメッセージの場合、この配列はnullです。 |
| afterImages | DML実行後のその行のデータ。DELETEメッセージの場合、この配列はnullです。 |
| sql | DDLのSQL文。 |
| tags | その他の情報。現在は空です。 |
| total | メッセージがシャーディングされている場合、シャーディング数を記録します。 |
| index | メッセージがシャーディングされている場合、現在のシャードのインデックスを記録します。 |
| beforeImageBytes | メッセージがシャーディングされている場合、シャーディング後のbeforeImages部分のデータを記録します。 |
| afterImageBytes | メッセージがシャーディングされている場合、シャーディング後のafterImages部分のデータを記録します。 |
Canal JSONメッセージ形式
{
"data": [
{
"id": "111",
"name": "name",
"description": "Big 2-wheel scooter",
"weight": "5.18"
}
],
"database": "inventory",
"es": 1589373560000,
"id": 9,
"isDdl": false,
"mysqlType": {
"id": "INTEGER",
"name": "VARCHAR(255)",
"description": "VARCHAR(512)",
"weight": "FLOAT"
},
"old": [
{
"weight": "5.15"
}
],
"pkNames": [
"id"
],
"sql": "",
"sqlType": {
"id": 4,
"name": 12,
"description": 12,
"weight": 7
},
"table": "products",
"ts": 1589373560798,
"type": "UPDATE"
}
パラメータ |
説明 |
|---|---|
| data | データ変更後の新しいレコードリスト。通常は配列で、各オブジェクトが1つのレコードを表し、フィールド名と値が含まれます。 |
| database | 変更が発生したデータベース名。 |
| es | イベントのタイムスタンプ。単位はミリ秒で、増分メッセージの実行時間を示します。 |
| id | 解釈されたメッセージID、連番。 |
| isDdl | DDL操作かどうか。true(テーブル構造変更などのDDL操作を示す)とfalse(データ変更を示す)が指定可能です。 |
| mysqlType | テーブルの各フィールドのMySQLタイプ。フィールド名をKey、タイプをValueとします。 |
| old | 変更前の古い値を記録します(変更されたフィールドとその古い値のみ含まれます)。 |
| pkNames | テーブルの主キーフィールド名のリスト。 |
| sql | DDLの場合、ここに対応するSQL文を記述します。DMLの場合、ここは空です。 |
| sqlType | テーブルの各フィールドのJava SQL Type。フィールド名をKey、タイプをValueとします。 |
| table | テーブル名。 |
| ts | メッセージを解釈したサーバーのタイムスタンプ、単位はミリ秒。 |
| type | 変更タイプ。INSERT、UPDATE、DELETE、またはDDLタイプが含まれます。 |
データベースフィールドマッピングの説明
OceanBaseデータベースMySQL互換モード
次の表は、OceanBaseデータベースのMySQL互換モードにおけるフィールド型とAvroプロトコルで定義されているデータ型とのマッピング関係を示しています。
分類 OceanBaseデータベースMySQL互換モードの型 対応するAvroの型 整数型 BOOL/BOOLEAN/TINYINT INTEGER SMALLINT INTEGER MEDIUMINT INTEGER INT/INTEGER INTEGER/LONG BIGINT LONG 固定小数点数型 DECIMAL DECIMAL NUMERIC DECIMAL 浮動小数点数型 FLOAT FLOAT DOUBLE DOUBLE ビット値型 BIT BIT 日付時刻型 DATETIME DATETIME TIMESTAMP TIMESTAMP DATE DATE TIME TIME YEAR INTEGER 文字列型 CHAR STRING VARCHAR STRING BINARY BINARY VARBINARY BINARY BLOB型 TINYBLOB BINARY BLOB BINARY MEDIUMBLOB BINARY LONGBLOB BINARY テキスト型 TINYTEXT STRING TEXT STRING MEDIUMTEXT STRING LONGTEXT STRING STRING STRING 列挙型 ENUM ENUM 集合型 SET SET JSONデータ型 JSON STRING 空間データ型 GEOMETRY GEOMETRY POINT GEOMETRY LINESTRING GEOMETRY POLYGON GEOMETRY MULTIPOINT GEOMETRY MULTILINESTRING GEOMETRY MULTIPOLYGON GEOMETRY GEOMETRYCOLLECTION GEOMETRY 高効率圧縮ビットマップデータ型 ROARINGBITMAP ROARINGBITMAP 配列データ型 ARRAY ARRAY マッピングデータ型 MAP MAP ベクトルデータ型 VECTOR VECTOR SPARSEVECTOR SPARSEVECTOR OceanBaseデータベースのOracle互換モード
次の表は、OceanBaseデータベースのOracle互換モードにおけるフィールド型とAvroプロトコルで定義されているデータ型とのマッピング関係を示しています。
分類 OceanBaseデータベースのOracle互換モードタイプ 対応するAvroの型 文字列型 CHAR STRING NCHAR STRING VARCHAR2 STRING VARCHAR STRING NVARCHAR2 STRING 数値型 NUMBER DECIMAL FLOAT FLOAT BINARY_FLOAT FLOAT BINARY_DOUBLE DOUBLE 日付時刻型 DATE DATE TIMESTAMP TIMESTAMP TIMESTAMP WITH TIME ZONE TIMESTAMP TIMESTAMP WITH LOCAL TIME ZONE TIMESTAMP 区間タイプ INTERVAL YEAR TO MONTH STRING INTERVAL DAY TO SECOND STRING ROWタイプ RAW RAW Rowidタイプ ROWID BINARY UROWID BINARY BLOB/CLOBタイプ BLOB BINARY CLOB STRING JSONタイプ JSON STRING XMLタイプ XML STRING 空間型 SDO_GEOMETRY GEOMETRY