- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
Ogg(Oracle GoldenGate)Format 是 Apache Flink 提供的一种 Changelog-Data-Capture(CDC)格式,它允许 Flink SQL 将 Ogg 捕获并同步到 Kafka 等消息系统的 JSON 变更事件,直接解析为INSERT/UPDATE/DELETE增量消息,也可以反向把 Flink SQL 中的变更消息编码为 Ogg JSON 输出到外部系统。读完本文,你将掌握 Ogg JSON 事件的结构与语义、如何通过 DDL 在 Kafka 上消费 Ogg 变更流、如何读取table/primary-keys等格式元数据,以及全部ogg-json.*配置项的取值与源码级实现原理。
Ogg Format 是什么
Oracle GoldenGate(简称 Ogg)是一个实时数据复制平台,通过数据库日志复制技术保证数据高可用并支撑实时分析。Ogg 为变更日志(changelog)提供了统一的格式 schema,并使用 JSON 完成消息序列化。Flink 的ogg-json格式正是针对这种 Ogg JSON 消息实现的序列化/反序列化 schema。
Flink 支持把 Ogg JSON 解释为 Flink SQL 系统中的INSERT/UPDATE/DELETE消息,典型应用场景包括:
- 将数据库的增量数据同步到其他系统;
- 审计日志处理;
- 基于数据库构建实时物化视图;
- 对数据库表的历史变化做 temporal join(时态关联)等。
同时,Flink 也支持把 Flink SQL 中的INSERT/UPDATE/DELETE消息编码为 Ogg JSON 并输出到 Kafka 等外部系统。需要特别注意的是:当前 Flink 无法把UPDATE_BEFORE和UPDATE_AFTER合并为一条UPDATE消息,因此在编码时 Flink 会把UPDATE_BEFORE编码为 Ogg 的 DELETE 消息、把UPDATE_AFTER编码为 Ogg 的 INSERT 消息(详见后文序列化源码分析)。
依赖引入
Ogg Json 依赖
ogg-json格式由flink-json模块提供,用户只需引入对应的 SQL jar 即可(flink-formats/flink-json模块的pom.xml对应 artifact)。该格式的工厂通过 META-INF/services/org.apache.flink.table.factories.Factory 注册,标识符为ogg-json。
提示:关于如何配置 Ogg Kafka Handler 把数据库变更日志同步到 Kafka topic,请参考 Ogg 官方 Kafka Handler 文档(例如 19.1 版本的Using the Kafka Handler)。
如何消费 Ogg 格式的数据
Ogg 为变更日志提供了统一格式,下面是一个从 OraclePRODUCTS表捕获到的更新(update)操作的 JSON 示例:
{ "before": { "id": 111, "name": "scooter", "description": "Big 2-wheel scooter", "weight": 5.18 }, "after": { "id": 111, "name": "scooter", "description": "Big 2-wheel scooter", "weight": 5.15 }, "op_type": "U", "op_ts": "2020-05-13 15:40:06.000000", "current_ts": "2020-05-13 15:40:07.000000", "primary_keys": [ "id" ], "pos": "00000000000000000000143", "table": "PRODUCTS" }提示:关于
before/after/op_type/op_ts/current_ts/primary_keys/pos/table等字段的具体含义,可以参考 Debezium 官方文档中 Oracle 连接器的事件字段说明(两者捕获字段语义高度一致)。
上述 OraclePRODUCTS表有 4 列(id、name、description、weight)。上面的 JSON 是一条针对该表的更新事件:id = 111这行数据的weight从5.18变更为5.15。假设这条消息已被同步到 Kafka topicproducts_ogg,可以使用如下 DDL 消费该 topic 并把变更事件解释为 changelog:
CREATE TABLE topic_products ( -- schema 与 Oracle "products" 表完全一致 id BIGINT, name STRING, description STRING, weight DECIMAL(10, 2) ) WITH ( 'connector' = 'kafka', 'topic' = 'products_ogg', 'properties.bootstrap.servers' = 'localhost:9092', 'properties.group.id' = 'testGroup', 'format' = 'ogg-json' );把 topic 注册为 Flink 表之后,即可把 Ogg 消息当作 changelog 数据源使用:
-- 在 Oracle "PRODUCTS" 表上构建实时物化视图 -- 计算同一产品 name 的最新平均 weight SELECT name, AVG(weight) FROM topic_products GROUP BY name; -- 把 Oracle "PRODUCTS" 表的全量数据及增量变更同步到 -- Elasticsearch 的 "products" 索引,用于后续搜索 INSERT INTO elasticsearch_products SELECT * FROM topic_products;反序列化时的 RowKind 映射
从源码 OggJsonDeserializationSchema 可以看到,Ogg 的op_type字段取值与 Flink 内部RowKind的对应关系为:
op_type取值 | 含义 | 映射为 Flink RowKind |
|---|---|---|
I | insert | INSERT,取after字段 |
U | update | UPDATE_BEFORE(before字段)+UPDATE_AFTER(after字段)两条消息 |
D | delete | DELETE,取before字段 |
T | truncate | 其他未知值时若未开启忽略解析错误则抛出异常 |
对于U(update)与D(delete)操作,如果before字段为 null,反序列化会抛出IllegalStateException。源码中给出的排查提示是:如果使用 Ogg Postgres Connector,需要确认 Postgres 表已设置REPLICA IDENTITY为FULL级别,否则无法拿到变更前的镜像数据(OggJsonDeserializationSchema)。另外,反序列化遇到 null 或空字节数组(Kafka 的 tombstone 消息)时会直接跳过,不会产出任何记录。
可用元数据(Available Metadata)
ogg-json格式可以将以下格式元数据暴露为表定义中的只读(VIRTUAL)列:
| Key | 数据类型 | 描述 |
|---|---|---|
table | STRING NULL | 完全限定的表名,格式为:目录名.模式名.表名(CATALOG NAME.SCHEMA NAME.TABLE NAME) |
primary-keys | ARRAY<STRING> NULL | 源表主键列名组成的数组;仅当 Ogg 侧配置属性includePrimaryKeys为 true 时,该字段才会出现在 JSON 输出中 |
ingestion-timestamp | TIMESTAMP_LTZ(6) NULL | 连接器处理该事件的时间戳,对应 Ogg 记录中的current_ts字段 |
event-timestamp | TIMESTAMP_LTZ(6) NULL | 源系统创建该事件的时间戳,对应 Ogg 记录中的op_ts字段 |
注意:格式元数据字段只有在对应 connector 转发格式元数据时才可用。目前只有 Kafka connector 能够为其 value format 暴露元数据字段。
从源码 OggJsonDecodingFormat.ReadableMetadata 可以看到上述元数据与 Ogg JSON 顶层字段的对应关系:table对应 JSON 顶层table,primary-keys对应 JSON 顶层primary_keys,ingestion-timestamp对应顶层current_ts(按yyyy-MM-dd'T'HH:mm:ss.SSSSSS格式解析),event-timestamp对应顶层op_ts(按yyyy-MM-dd HH:mm:ss.SSSSSS格式解析)。
下面的示例展示了如何在 Kafka 表上访问 Ogg 元数据字段:
CREATE TABLE KafkaTable ( origin_ts TIMESTAMP(3) METADATA FROM 'value.ingestion-timestamp' VIRTUAL, event_time TIMESTAMP(3) METADATA FROM 'value.event-timestamp' VIRTUAL, origin_table STRING METADATA FROM 'value.table' VIRTUAL, primary_keys ARRAY<STRING> METADATA FROM 'value.primary-keys' VIRTUAL, user_id BIGINT, item_id BIGINT, behavior STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'user_behavior', 'properties.bootstrap.servers' = 'localhost:9092', 'properties.group.id' = 'testGroup', 'scan.startup.mode' = 'earliest-offset', 'value.format' = 'ogg-json' );格式选项(Format Options)
ogg-json格式支持以下配置选项,均定义于 OggJsonFormatFactory 与 OggJsonFormatOptions 中:
| 选项 | 是否必填 | 默认值 | 类型 | 描述 |
|---|---|---|---|---|
format | 必填 | (无) | String | 指定使用的格式,此处应为'ogg-json' |
ogg-json.ignore-parse-errors | 可选 | false | Boolean | 遇到解析错误时跳过对应字段和行而不是失败;出错时字段会被置为 null |
ogg-json.timestamp-format.standard | 可选 | 'SQL' | String | 指定输入/输出的时间戳格式,目前支持'SQL'与'ISO-8601'两种取值(详见下方说明) |
ogg-json.map-null-key.mode | 可选 | 'FAIL' | String | 序列化 Map 数据时遇到 null key 的处理模式,支持'FAIL'、'DROP'、'LITERAL' |
ogg-json.map-null-key.literal | 可选 | 'null' | String | 当ogg-json.map-null-key.mode为LITERAL时,用于替换 null key 的字符串字面量 |
ogg-json.encode.ignore-null-fields | 可选 | false | Boolean | 只编码非 null 字段;默认会包含所有字段 |
ogg-json.timestamp-format.standard两种取值的差异:
'SQL':按yyyy-MM-dd HH:mm:ss.s{precision}格式解析输入时间戳(例如2020-12-30 12:13:14.123),输出也采用同样格式;'ISO-8601':按yyyy-MM-ddTHH:mm:ss.s{precision}格式解析输入时间戳(例如2020-12-30T12:13:14.123),输出也采用同样格式。
ogg-json.map-null-key.mode三种取值的差异:
'FAIL':遇到 Map 的 null key 时抛出异常;'DROP':丢弃 Map 数据中 null key 的条目;'LITERAL':用字符串字面量替换 null key,字面量由ogg-json.map-null-key.literal选项指定。
工厂测试 OggJsonFormatFactoryTest 对这些选项的取值校验给出了明确证据:ogg-json.ignore-parse-errors只接受布尔值(true/false,不区分大小写);ogg-json.timestamp-format.standard仅支持SQL和ISO-8601;ogg-json.map-null-key.mode仅支持LITERAL、FAIL、DROP,传入非法值会抛出ValidationException。
此外,从工厂源码还可以看到,编码侧还支持从通用 JSON 格式继承的ogg-json.encode.decimal-as-plain-number选项(将 DECIMAL 编码为普通数字而非字符串),它并非ogg-json的专属选项但同样生效。
数据类型的映射(Data Type Mapping)
当前 Ogg 格式使用 JSON 格式完成序列化与反序列化,因此其数据类型映射规则与 Flink 的 JSON Format 完全一致,包括时间戳精度处理、DECIMAL的编码形式(字符串或普通数字)、ARRAY/MAP/ROW等复合类型的映射方式等,可直接参考 JSON Format 文档中的 Data Type Mapping 一节。
源码级实现原理:解码与编码链路
解码链路:从 Ogg JSON 到 RowData
ogg-json的反序列化由OggJsonDeserializationSchema完成(flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/ogg/OggJsonDeserializationSchema.java)。其内部构造的 JSON 行类型固定为:
ROW( "before" <物理数据类型>, "after" <物理数据类型>, "op_type" STRING )即反序列化时只关心before、after、op_type三个顶层字段;table、primary_keys、current_ts、op_ts等字段只有在声明了对应元数据列时才会被追加到根 RowType 中用于元数据提取。解码完成后按上文表格中的op_type规则设置RowKind并产出记录;对于 update 操作会先后产出UPDATE_BEFORE与UPDATE_AFTER两条消息。解码格式声明其 changelog 模式同时包含INSERT、UPDATE_BEFORE、UPDATE_AFTER、DELETE四种 RowKind(见 OggJsonDecodingFormat#getChangelogMode)。
编码链路:从 RowData 到 Ogg JSON
ogg-json的序列化由OggJsonSerializationSchema完成(flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/ogg/OggJsonSerializationSchema.java)。序列化时同样只输出before、after、op_type三个顶层字段(源码注释明确说明 Ogg JSON 中的source、ts_ms等其他信息在此处并不需要):
INSERT/UPDATE_AFTER:before置为 null,after写入当前行数据,op_type置为I;UPDATE_BEFORE/DELETE:before写入当前行数据,after置为 null,op_type置为D。
这正是文档中所说"Flink 把 UPDATE_BEFORE / UPDATE_AFTER 分别编码为 DELETE / INSERT 两条 Ogg 消息"的底层实现。编码格式的 changelog 模式同样声明支持四种 RowKind(见 OggJsonFormatFactory)。
测试验证与实战建议
仓库在 flink-formats/flink-json/src/test/java/org/apache/flink/formats/json/ogg/ 目录下提供了完整的验证用例:
- OggJsonSerDeSchemaTest:基于测试资源
ogg-data.txt覆盖了完整的 INSERT / UPDATE / DELETE 序列化与反序列化往返,并验证了元数据列(table、primary-keys、ingestion-timestamp、event-timestamp)的读取结果,以及 null / 空字节 tombstone 消息被跳过、不产生任何记录的行为; - OggJsonFormatFactoryTest:验证工厂对全部选项的解析与非法值校验;
- OggJsonFileSystemITCase:验证文件系统连接器场景下 Ogg JSON 的端到端读写。
实战中的几点建议:
- 主键与镜像数据:使用 Ogg Postgres Connector 时,务必把源表
REPLICA IDENTITY设置为FULL,否则 update/delete 事件缺少before数据会导致消费失败; - 元数据声明:Kafka 消费场景下,
table、primary-keys、ingestion-timestamp、event-timestamp必须声明为METADATA FROM 'value.xxx' VIRTUAL才能读取; - 时间戳格式:Ogg 消息中
op_ts/current_ts默认形如2020-05-13 15:40:06.000000(含空格),与ogg-json.timestamp-format.standard的'SQL'默认解析格式匹配;若 Ogg 侧配置输出 ISO-8601 风格(T分隔符),则需显式设置'ogg-json.timestamp-format.standard' = 'ISO-8601'; - 编码语义:当使用
ogg-json作为 sink 格式时,上游的 update 会以"先 DELETE(before)后 INSERT(after)"两条 Ogg 消息落盘,下游消费者需要按此语义还原变更。
- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
相关推荐
Flink Ogg Format 深度指南:Oracle GoldenGate 变更日志的实时接入与输出
Flink Ogg Format 深度指南:Oracle GoldenGate 变更日志的实时接入与输出 Oracle GoldenGate(简称 Ogg)是
大数据流处理批处理数据工程Flink Canal Format 实战指南:基于 canal-json 的 MySQL CDC 变更数据捕获与同步
Flink Canal Format 实战指南:基于 canal json 的 MySQL CDC 变更数据捕获与同步 Canal 是阿里巴巴开源的 CDC(C
大数据流处理批处理数据工程Flink 集成 Maxwell JSON 格式:基于 MySQL CDC 的 Changelog 流接入与输出完整指南
Flink 集成 Maxwell JSON 格式:基于 MySQL CDC 的 Changelog 流接入与输出完整指南 Maxwell 是业界常用的 CDC(
大数据流处理批处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考