本記事では、データサブスクリプションでサポートされているデータ型である 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-Value型 BIT BIT 日付時刻型 DATETIME DATETIME TIMESTAMP TIMESTAMP DATE DATE TIME TIME YEAR INTEGER 文字列型 CHAR STRING VARCHAR STRING BINARY BINARY VARBINARY BINARY LOB型 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 LOB型 BLOB BINARY CLOB STRING JSON型 JSON STRING XML型 XML STRING 空間型 SDO_GEOMETRY GEOMETRY