Flink Ogg Format 实战:基于 Oracle GoldenGate JSON 的 Changelog 数据接入指南
2026/9/23 14:36:40 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/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_BEFOREUPDATE_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 列(idnamedescriptionweight)。上面的 JSON 是一条针对该表的更新事件:id = 111这行数据的weight5.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
IinsertINSERT,取after字段
UupdateUPDATE_BEFOREbefore字段)+UPDATE_AFTERafter字段)两条消息
DdeleteDELETE,取before字段
Ttruncate其他未知值时若未开启忽略解析错误则抛出异常

对于U(update)与D(delete)操作,如果before字段为 null,反序列化会抛出IllegalStateException。源码中给出的排查提示是:如果使用 Ogg Postgres Connector,需要确认 Postgres 表已设置REPLICA IDENTITYFULL级别,否则无法拿到变更前的镜像数据(OggJsonDeserializationSchema)。另外,反序列化遇到 null 或空字节数组(Kafka 的 tombstone 消息)时会直接跳过,不会产出任何记录。

可用元数据(Available Metadata)

ogg-json格式可以将以下格式元数据暴露为表定义中的只读(VIRTUAL)列:

Key数据类型描述
tableSTRING NULL完全限定的表名,格式为:目录名.模式名.表名(CATALOG NAME.SCHEMA NAME.TABLE NAME)
primary-keysARRAY<STRING> NULL源表主键列名组成的数组;仅当 Ogg 侧配置属性includePrimaryKeys为 true 时,该字段才会出现在 JSON 输出中
ingestion-timestampTIMESTAMP_LTZ(6) NULL连接器处理该事件的时间戳,对应 Ogg 记录中的current_ts字段
event-timestampTIMESTAMP_LTZ(6) NULL源系统创建该事件的时间戳,对应 Ogg 记录中的op_ts字段

注意:格式元数据字段只有在对应 connector 转发格式元数据时才可用。目前只有 Kafka connector 能够为其 value format 暴露元数据字段。

从源码 OggJsonDecodingFormat.ReadableMetadata 可以看到上述元数据与 Ogg JSON 顶层字段的对应关系:table对应 JSON 顶层tableprimary-keys对应 JSON 顶层primary_keysingestion-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可选falseBoolean遇到解析错误时跳过对应字段和行而不是失败;出错时字段会被置为 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'Stringogg-json.map-null-key.modeLITERAL时,用于替换 null key 的字符串字面量
ogg-json.encode.ignore-null-fields可选falseBoolean只编码非 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仅支持SQLISO-8601ogg-json.map-null-key.mode仅支持LITERALFAILDROP,传入非法值会抛出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 )

即反序列化时只关心beforeafterop_type三个顶层字段;tableprimary_keyscurrent_tsop_ts等字段只有在声明了对应元数据列时才会被追加到根 RowType 中用于元数据提取。解码完成后按上文表格中的op_type规则设置RowKind并产出记录;对于 update 操作会先后产出UPDATE_BEFOREUPDATE_AFTER两条消息。解码格式声明其 changelog 模式同时包含INSERTUPDATE_BEFOREUPDATE_AFTERDELETE四种 RowKind(见 OggJsonDecodingFormat#getChangelogMode)。

编码链路:从 RowData 到 Ogg JSON

ogg-json的序列化由OggJsonSerializationSchema完成(flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/ogg/OggJsonSerializationSchema.java)。序列化时同样只输出beforeafterop_type三个顶层字段(源码注释明确说明 Ogg JSON 中的sourcets_ms等其他信息在此处并不需要):

  • INSERT/UPDATE_AFTERbefore置为 null,after写入当前行数据,op_type置为I
  • UPDATE_BEFORE/DELETEbefore写入当前行数据,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 序列化与反序列化往返,并验证了元数据列(tableprimary-keysingestion-timestampevent-timestamp)的读取结果,以及 null / 空字节 tombstone 消息被跳过、不产生任何记录的行为;
  • OggJsonFormatFactoryTest:验证工厂对全部选项的解析与非法值校验;
  • OggJsonFileSystemITCase:验证文件系统连接器场景下 Ogg JSON 的端到端读写。

实战中的几点建议:

  1. 主键与镜像数据:使用 Ogg Postgres Connector 时,务必把源表REPLICA IDENTITY设置为FULL,否则 update/delete 事件缺少before数据会导致消费失败;
  2. 元数据声明:Kafka 消费场景下,tableprimary-keysingestion-timestampevent-timestamp必须声明为METADATA FROM 'value.xxx' VIRTUAL才能读取;
  3. 时间戳格式: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'
  4. 编码语义:当使用ogg-json作为 sink 格式时,上游的 update 会以"先 DELETE(before)后 INSERT(after)"两条 Ogg 消息落盘,下游消费者需要按此语义还原变更。
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询