SeaTunnel OGG_JSON 格式详解:以 Oracle GoldenGate 统一日志格式接入 CDC 实时同步

发布时间:2026/9/18 16:29:49
SeaTunnel OGG_JSON 格式详解:以 Oracle GoldenGate 统一日志格式接入 CDC 实时同步
SeaTunnel OGG_JSON 格式详解以 Oracle GoldenGate 统一日志格式接入 CDC 实时同步【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文基于 SeaTunnel 官方文档 docs/zh/connectors/formats/ogg-json.md 展开讲解ogg_json格式如何把 Oracle GoldenGateOGG的 JSON 变更日志解释为 SeaTunnel 内部的 INSERT/UPDATE/DELETE 消息以及如何把 SeaTunnel 行消息反向序列化回 Ogg JSON。读完后你将掌握Ogg JSON 消息的字段结构与事件类型映射规则、ogg_json.*各格式选项的语义与底层实现、通过 Kafka 消费 Ogg 变更日志并写入 MySQL 的完整作业配置以及 UPDATE 事件被拆分为 DELETE INSERT 的机制与源码级依据。一、OGG 格式是什么SeaTunnel 如何使用它Oracle GoldenGatea.k.a OGG是一项托管服务提供实时数据网格平台它使用复制来保持数据高度可用并支持实时分析客户可以设计、执行和监控数据复制和流处理方案而无需分配或管理计算环境。对 SeaTunnel 而言OGG 的价值在于它为变更日志changelog提供了一套统一的格式结构并支持使用 JSON 序列化消息——这使得不同数据库、不同复制链路上产生的变更事件可以用同一种 schema 表达下游系统只需实现一套解析逻辑。SeaTunnel 对 Ogg JSON 的支持体现在两个方向反序列化Source 侧将 Ogg JSON 消息解释为 SeaTunnel 系统中的INSERT/UPDATE/DELETE消息。典型场景包括将增量数据从数据库同步到其他系统审计日志数据库的实时物化视图关联维度数据库的变更历史等等。序列化Sink 侧将 SeaTunnel 中的INSERT/UPDATE/DELETE消息转化为 Ogg JSON 消息并发送到类似 Kafka 这样的存储中。需要特别注意的一个限制目前 SeaTunnel 无法将UPDATE_BEFORE和UPDATE_AFTER组合成单个 UPDATE 消息因此默认情况下会将二者转化为一条 DELETE 和一条 INSERT 的 Ogg 消息来实现。这一限制的底层实现与变通方式见本文第五节。格式实现代码位于seatunnel-formats/seatunnel-format-json模块的ogg包下OggJsonDeserializationSchema.java反序列化Ogg JSON → SeaTunnelRowOggJsonSerializationSchema.java序列化SeaTunnelRow → Ogg JSONOggJsonFormatOptions.java格式选项定义OggJsonSerDeSchemaTest.java序列化/反序列化的完整测试用例二、Ogg JSON 变更消息的字段结构以从 OraclePRODUCTS表捕获的一条更新操作为例这是从 Oracle 复制链路同步出来的典型消息{ 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 }各字段含义可参考 Debezium 对同类 change event 的定义Ogg JSON 与之高度兼容。上面这条消息对应PRODUCTS表4 列id、name、description、weight上的一次更新事件id 111的行的weight字段从 5.18 改为 5.15。从 OggJsonDeserializationSchema.java 的常量定义看SeaTunnel 实际解析时只依赖以下字段字段用途op_type操作类型取值IINSERT/UUPDATE/DDELETEbefore/after变更前/后的行数据按用户声明的 schema 解析op_ts变更发生时间格式为2020-05-13 15:40:07.000000这样的日期字符串解析后转换为 Unix 毫秒时间戳写入行的 event_time 元数据table表名可含database.table形式用于库表过滤current_ts、primary_keys、pos等字段解析时会被忽略——源码中createJsonRowType方法的注释也明确写道Ogg JSON contains other information, e.g. ts, sql, but we dont need them。三、格式选项说明在连接器配置中使用该格式时需要声明format ogg_json并可选配置以下选项选项默认值是否需要描述format(none)是指定要使用的格式这里应该是ogg_jsonogg_json.ignore-parse-errorsfalse否跳过有解析错误的字段和行而不是失败。如果出现错误字段将设置为 nullogg_json.database.include(none)否正则表达式可选。通过正则匹配 Ogg 记录中的database元字段仅读取特定数据库的变更日志行。Pattern 与 Java 的Pattern兼容ogg_json.table.include(none)否正则表达式可选。通过正则匹配 Ogg 记录中的table元字段仅读取特定表的更改日志行。Pattern 与 Java 的Pattern兼容对应源码OggJsonFormatOptions.java 中定义了DATABASE_INCLUDEkey 为database.include和TABLE_INCLUDEkey 为table.include两个Option二者均为无默认值的字符串类型且描述明确声明 Pattern 与 JavaPattern兼容IGNORE_PARSE_ERRORS复用 JSON 格式共用的同名选项读取时以Boolean.parseBoolean解析缺省为false。一个值得注意的细节在 Kafka Source 的当前实现中ogg_json格式的反序列化 schema 被构造时固定开启了忽略解析错误。见 KafkaSourceConfig.javacase OGG_JSON: schema OggJsonDeserializationSchema.builder(catalogTable) .setIgnoreParseErrors(true) .build(); break;也就是说从源码结构看经由 Kafka 连接器消费 Ogg JSON 时遇到脏数据JSON 解析失败或消息内容异常会被静默跳过而不是让作业失败。如果你在其它连接器中直接使用格式 API则ogg_json.ignore-parse-errors的文档默认值false仍然生效行为上会抛出CommonError.jsonOperationError。四、反序列化原理Ogg JSON 到 SeaTunnelRow 的事件映射核心逻辑在 OggJsonDeserializationSchema.java 的deserializeMessage方法中可以归纳为以下处理流程Tombstone 消息跳过message为 null 或长度为 0 时直接返回不做任何处理这正好配合 Kafka 中 delete tombstone 的语义。库表过滤若配置了database.include/table.include取消息中table字段按.拆分——第 0 段匹配库正则、第 1 段匹配表正则任一不匹配则丢弃该条消息。注意这里使用的是matcher(...).matches()即全量匹配而非部分匹配例如测试用例 OggJsonSerDeSchemaTest.java 中使用^OG.*与^TBL.*过滤对应测试数据文件 ogg-data-filter-table.txt。按op_type分派事件IINSERT取after数据转为 SeaTunnelRow行类型为INSERTUUPDATE分别取before与after转为两行行类型为UPDATE_BEFORE/UPDATE_AFTER一次消息产出两条行DDELETE取before数据转为 SeaTunnelRow行类型为DELETE其它取值抛出IllegalStateException(Unknown operation type ...)。时间戳处理若消息含op_ts用DateTimeUtils.parse解析为 UTC 纪元毫秒后通过MetadataUtil.setEventTime写入行的 event_time 元数据供下游做事件时间语义。表路径标注若 CatalogTable 携带表路径会写入每行的tableIdsetTableId序列化时再映射回table字段。错误处理before字段为 null 的 UPDATE/DELETE 消息会抛出带有明确提示的异常提示文案指向 Postgres 场景The before field of %s operation message is null, if you are using Ogg Postgres Connector, please check the Postgres table has been set REPLICA IDENTITY to FULL level.这是一个非常实用的排障指引当使用 OGG Postgres Connector 时UPDATE/DELETE 事件缺失 before 镜像的根因通常是表没有将REPLICA IDENTITY设置为FULL级别此时应到源库执行ALTER TABLE xxx REPLICA IDENTITY FULL;。该行为同样有测试覆盖OggJsonSerDeSchemaTest.java 中testDeserializeNoDataJson用{op_type:U}断言了这条异常信息。若配置了ignore-parse-errors上述所有RuntimeExceptionJSON 解析失败、before 缺失、未知 op_type 等都会被吞掉坏行被跳过测试用例testDeserializeNoJson输入{]与testDeserializeEmptyJson输入{}验证了未开启忽略时抛出CommonError.jsonOperationError的行为。五、序列化原理SeaTunnelRow 到 Ogg JSONSink 方向由 OggJsonSerializationSchema.java 实现。它把一行数据包装成固定 5 列的中间结构后交给通用 JSON 序列化器输出// createJsonRowType: 输出行的列结构 new String[] {before, after, op_type, table, op_ts}行类型RowKind到op_type的映射规则在rowKind2String方法中OggJsonSerializationSchema.javaSeaTunnel 行类型序列化结果INSERTop_type I数据放入afterbefore为 nullUPDATE_AFTER默认模式拆分为一条D消息携带 before 镜像 一条I消息UPDATE_AFTERmergeUpdateEventFlag true时合并为一条op_type U消息before/after同时填充UPDATE_BEFORE默认模式下映射为op_type DDELETEop_type D数据放入beforeafter为 null这正对应文档开头所述的限制目前 SeaTunnel 无法将 UPDATE_BEFORE 和 UPDATE_AFTER 组合成单个 UPDATE 消息……将 UPDATE_BEFORE 和 UPDATE_AFTER 转化为 DELETE 和 INSERT Ogg 消息来实现。从源码结构看类中已经提供了mergeUpdateEventFlag构造参数mergeUpdateEventFlag true时缓存UPDATE_BEFORE行、待UPDATE_AFTER到来时合并输出U消息before字段即被填充说明合并能力已在序列化器层面具备但在 Kafka 连接器的 Sink 侧DefaultSeaTunnelRowSerializer.java 目前构造的是无参版本new OggJsonSerializationSchema(rowType)即走拆分为 DELETE INSERT的默认路径。测试用例 OggJsonSerDeSchemaTest.java 完整验证了两种模式的输出例如默认模式下一次 UPDATE 会产出{before:null,after:{id:106,name:hammer,description:16oz carpenters hammer,weight:1.0},op_type:D,table:..test,op_ts:1589390787000} {before:null,after:{id:106,name:hammer,description:18oz carpenter hammer,weight:1.0},op_type:I,table:..test,op_ts:1589390787000}而开启合并标志后则产出单条U消息{before:{id:106,name:hammer,description:16oz carpenters hammer,weight:1.0},after:{id:106,name:hammer,description:18oz carpenter hammer,weight:1.0},op_type:U,table:..test,op_ts:1589390787000}此外serialize方法还会把行的tableId写入table字段、把 event_time 元数据写入op_tsLONG 毫秒从而与上游解析出的字段保持闭环。六、实战Kafka 消费 Ogg 变更日志并写入 MySQL假设上面PRODUCTS表的变更消息已经通过复制链路同步到了 Kafka topic例如 topic 名为ogg下面的 SeaTunnel 作业可以消费该 topic 并把变更事件落库到 MySQL注意该作业声明为 STREAMING 模式env { parallelism 1 job.mode STREAMING } source { Kafka { bootstrap.servers 127.0.0.1:9092 topic ogg plugin_output kafka_name start_mode earliest schema { fields { id int name string description string weight double } }, format ogg_json } } sink { jdbc { url jdbc:mysql://127.0.0.1/test driver com.mysql.cj.jdbc.Driver user root password 12345678 table ogg primary_keys [id] } }配置要点说明format ogg_json指定 Kafka 消息的解析格式。Kafka 连接器支持的消息格式由 MessageFormat.java 枚举定义其中包含OGG_JSONschema { fields { ... } }显式声明业务字段的名称与类型。由于 Ogg JSON 的before/after只是行数据的容器SeaTunnel 需要用户声明目标表结构反序列化器OggJsonDeserializationSchema内部的JsonDeserializationSchema按此 schema 逐字段取值转换由于op_type U的消息会展开为UPDATE_BEFOREUPDATE_AFTER两条行JDBC Sink 收到后会按主键primary_keys [id]先执行删除再插入最终库中数据与源端保持一致若 topic 中混入了不关心的库表可追加ogg_json.database.include/ogg_json.table.include正则进行过滤配合ogg_json.ignore-parse-errors控制脏数据行为注意前文提到的 Kafka Source 侧对ogg_json已默认开启忽略解析错误。七、测试用例与排错指引单元测试 OggJsonSerDeSchemaTest.java 覆盖了 Ogg 格式的全部关键行为遇到线上问题时可对照这些断言快速定位测试方法覆盖场景testFilteringTablesdatabase.include/table.include正则过滤数据文件ogg-data-filter-table.txttestDeserializeNullRownull/空消息tombstone被安全跳过testDeserializeNoJson/testDeserializeEmptyJson非法 JSON{]与空对象{}抛出jsonOperationErrortestDeserializeNoDataJsonUPDATE 消息缺before字段时抛出 REPLICA IDENTITY FULL 提示testDeserializeUnknownTypeJsonop_type XX时抛出Unknown operation typerunTestround-trip完整验证反序列化 → 行事件序列 → 再序列化的输出含默认拆分模式与merge_update_event合并模式的逐字节比对综合起来ogg_json格式的典型排错路径是作业报jsonOperationError且消息以{开头但字段不全 → 检查源端复制链路是否输出完整消息或确认 Kafka Source 是否按预期静默跳过了坏行UPDATE/DELETE 事件报 before 为 null → 源库表需要REPLICA IDENTITY FULLPostgres 场景库表中混入无关变更 → 使用ogg_json.database.include/ogg_json.table.include正则注意是全量匹配过滤下游期望单条 UPDATE 语义 → 当前经 Kafka Sink 时 UPDATE 会呈现为一对 DELETE INSERT 消息这是该版本的既定行为。小结ogg_json格式让 SeaTunnel 能够以低耦合的方式接入 Oracle GoldenGate 复制链路Source 侧将I/U/D事件映射为 SeaTunnel 的行变更事件U展开为 UPDATE_BEFORE UPDATE_AFTER 两行并附带 event_time 与表路径元数据Sink 侧将行事件序列化为统一的五字段 Ogg JSONbefore/after/op_type/table/op_ts。配合 Kafka 连接器与显式 schema 声明即可实现Ogg 变更日志 → 流式消费 → 目标数据库的端到端增量同步这也是增量数据同步、审计日志、实时物化视图等场景的通用底座。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考