- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
本文基于 SeaTunnel 官方中文文档(docs/zh/connector-v2/formats/ogg-json.md),结合seatunnel-formats/seatunnel-format-json模块中 Ogg JSON 反序列化/序列化源码与端到端测试配置,系统讲解 Ogg(Oracle GoldenGate)JSON 格式在 SeaTunnel 中的完整使用方式。读者将掌握:Ogg JSON 消息的结构与字段语义、ogg_json全部格式选项的配置方法、Kafka 中消费 Ogg 变更日志并同步到下游数据库的完整作业写法,以及 SeaTunnel 将 INSERT/UPDATE/DELETE 变更消息反向编码为 Ogg JSON 的底层机制与限制。
Ogg 与 Ogg JSON 格式概览
Oracle GoldenGate(简称 Ogg)是 Oracle 提供的一项基于复制技术的实时数据集成服务,通过数据库复制保持数据高可用并支撑实时分析,用户无需自行分配或管理计算环境即可设计、执行和监控数据复制与流数据处理方案。Ogg 为变更日志提供了统一的结构化格式,并支持使用 JSON 序列化消息。
SeaTunnel 对 Ogg JSON 的支持体现在两个方向:
- 解析(反序列化):将 Ogg JSON 消息解释为 SeaTunnel 内部的 INSERT / UPDATE / DELETE 变更消息,从而把 Ogg 捕获的数据库增量变更接入 SeaTunnel 作业;
- 编码(序列化):将 SeaTunnel 中的 INSERT / UPDATE / DELETE 变更消息转化为 Ogg JSON 消息,并发送到 Kafka 等存储。
这一能力对应着多个典型的实时数据应用场景:
- 将增量数据从数据库同步到其他系统;
- 构建审计日志;
- 实现数据库的实时物化视图;
- 关联维度数据库的变更历史等。
需要特别说明的是:SeaTunnel目前无法将 UPDATE_BEFORE 和 UPDATE_AFTER 组合成单个 UPDATE 消息,因此编码阶段会将 UPDATE_BEFORE 与 UPDATE_AFTER 分别转换为 DELETE 与 INSERT 两种 Ogg 消息来实现(详见下文"变更消息的序列化"章节)。
Ogg JSON 消息结构详解
Ogg 为变更日志提供了统一的消息格式。以下是一条从 OraclePRODUCTS表捕获的更新操作示例(该表包含id、name、description、weight四列):
{ "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" }上面这条 JSON 消息是products表上的一个更新变更事件:id = 111的行的weight字段值从5.18变更为5.15。各字段的完整含义可参考 Ogg 官方数据变更事件文档(原文档指向 Debezium Oracle connector 的数据变更事件章节)。
从 SeaTunnel 源码(OggJsonDeserializationSchema.java)看,解析阶段实际依赖的关键字段及其取值如下:
| 字段 | 取值 | 用途 |
|---|---|---|
op_type | I(INSERT)、U(UPDATE)、D(DELETE) | 决定消息被解释为插入、更新还是删除事件 |
before | 变更前的整行数据(JSON 对象) | UPDATE / DELETE 事件读取,用于产生UPDATE_BEFORE/DELETE行 |
after | 变更后的整行数据(JSON 对象) | INSERT / UPDATE 事件读取,用于产生INSERT/UPDATE_AFTER行 |
table | 形如database.table的元字段 | 供ogg_json.database.include/ogg_json.table.include正则过滤使用 |
op_ts、current_ts、primary_keys、pos等 | 时间戳、主键、位点等元信息 | 反序列化时会被忽略(源码注释明确说明 Ogg JSON 中的ts、sql等附加信息在 SeaTunnel 中不需要) |
其中table元字段的过滤逻辑较为特殊:源码在匹配数据库与表时,会将该字段按.分割为两段,第一段与database.include正则匹配,第二段与table.include正则匹配(OggJsonDeserializationSchema.java),因此记录中的table字段通常应形如库名.表名。
格式选项说明
使用ogg_json格式时,需要在连接器(如 Kafka Source / Sink)的配置中通过format = ogg_json指定格式,并可搭配以下选项(详见 OggJsonFormatOptions.java):
| 选项 | 默认值 | 是否必需 | 描述 |
|---|---|---|---|
format | (none) | 是 | 指定要使用的格式,这里应为ogg_json |
ogg_json.ignore-parse-errors | false | 否 | 跳过有解析错误的字段和行而不是失败;出现错误时字段会被设置为null |
ogg_json.database.include | (none) | 否 | 可选正则表达式,通过匹配 Ogg 记录中的database元字段来仅读取特定数据库的变更日志行;该字符串的 Pattern 模式与 Java 的Pattern兼容 |
ogg_json.table.include | (none) | 否 | 可选正则表达式,通过匹配 Ogg 记录中的table元字段来仅读取特定表的变更日志行;该字符串的 Pattern 模式与 Java 的Pattern兼容 |
从源码可以确认这些选项的解析方式:
ogg_json.ignore-parse-errors对应JsonFormatOptions.IGNORE_PARSE_ERRORS,默认值为false,即默认任何解析错误都会使作业失败;ogg_json.database.include与ogg_json.table.include均无默认值(noDefaultValue),未配置时表示不过滤任何数据库 / 表;- 两个 include 正则最终会被编译为 Java
Pattern(OggJsonDeserializationSchema.java),因此支持完整的 Java 正则语法,例如^OG.*、^TBL.*这类前缀匹配写法。
Kafka 消费 Ogg 变更日志实战
假设 OraclePRODUCTS表的 Ogg 变更消息已经同步到 Kafka topic(如ogg),可以使用下面的 SeaTunnel 作业配置来消费该 topic,将变更事件解析为 SeaTunnel 的行消息,并写入 MySQL 目标表:
env { parallelism = 1 job.mode = "STREAMING" } source { Kafka { bootstrap.servers = "127.0.0.1:9092" topic = "ogg" result_table_name = "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 Source 使用 Ogg JSON 反序列化器解析每条消息;schema.fields定义了业务数据的字段类型,需与 Ogg 捕获的源表列一致(本示例中id为int,name/description为string,weight为double),SeaTunnel 据此把before/after中的 JSON 对象转换为行数据;job.mode = "STREAMING"配合start_mode = earliest表示以流式模式从最早位点开始持续消费;- 反序列化时,
op_type为I的消息产生INSERT(+I)行,为U的消息同时产生UPDATE_BEFORE(-U)与UPDATE_AFTER(+U)两行,为D的消息产生DELETE(-D)行,下游 JDBC Sink 会根据行的 RowKind 执行对应的写入语义。
变更消息的反序列化语义
Ogg JSON 反序列化器位于 OggJsonDeserializationSchema.java,其核心逻辑(deserializeMessage方法)遵循以下流程:
- 跳过墓碑消息:如果消息为
null或长度为 0(Kafka 的 tombstone 记录),直接返回不产生任何输出行; - 解析 JSON:将字节流解析为 Jackson 的
ObjectNode,解析失败时若配置了ogg_json.ignore-parse-errors = true则静默跳过,否则抛出SeaTunnelRuntimeException; - 数据库 / 表过滤:若配置了 include 正则,按上文所述方式匹配
table元字段拆分出的库名与表名,不匹配的消息被丢弃; - 按操作类型分发:
I(INSERT):读取after节点转换为行,以RowKind.INSERT输出;U(UPDATE):读取before与after节点,分别以RowKind.UPDATE_BEFORE和RowKind.UPDATE_AFTER输出两行;若before为空则抛出IllegalStateException;D(DELETE):读取before节点以RowKind.DELETE输出,同样要求before非空;- 其他操作类型:抛出"Unknown operation type"异常。
值得注意的约束是:UPDATE 与 DELETE 事件要求before字段非空。源码中专门定义了REPLICA_IDENTITY_EXCEPTION错误信息(OggJsonDeserializationSchema.java),提示:如果使用 Ogg Postgres Connector 且 UPDATE / DELETE 消息的before字段为 null,需要检查 Postgres 表的REPLICA IDENTITY是否设置为FULL级别,否则无法提供变更前镜像数据。
单元测试 OggJsonSerDeSchemaTest.java 也验证了这些行为,例如:
testDeserializeNullRow:空消息不产生任何输出行;testDeserializeNoJson/testDeserializeEmptyJson:非法 JSON 直接抛出解析异常;testDeserializeNoDataJson:{"op_type":"U"}这类缺少before的消息抛出REPLICA IDENTITY相关异常;testDeserializeUnknownTypeJson:未知操作类型抛出Unknown operation type 'XX';testFilteringTables:通过setDatabase("^OG.*").setTable("^TBL.*")验证库表正则过滤;- 完整的序列化/反序列化往返断言:UPDATE 事件在反序列化时拆成
-U与+U两行,在序列化时又分别编码为{"type":"DELETE"}与{"type":"INSERT"}。
变更消息的序列化:将 SeaTunnel 变更流编码为 Ogg JSON
SeaTunnel 同样支持把内部的 INSERT/UPDATE/DELETE 变更消息编码为 Ogg JSON 输出到 Kafka 等存储,实现变更日志的"接力"转发。序列化器位于 OggJsonSerializationSchema.java,其输出结构与 Ogg 标准消息有所差异,采用如下简化形态:
{"data":{"id":111,"name":"scooter","description":"Big 2-wheel scooter","weight":5.15},"type":"INSERT"}输出对象只包含两个字段:data(整行数据,类型为源 schema)与type(操作类型字符串)。RowKind 到type的映射规则由rowKind2String方法决定:
| SeaTunnel RowKind | 输出的 Oggtype |
|---|---|
INSERT | INSERT |
UPDATE_AFTER | INSERT |
UPDATE_BEFORE | DELETE |
DELETE | DELETE |
| 其他 | 抛出UNSUPPORTED_OPERATION异常 |
这正是原文档所述限制的具体体现:由于 SeaTunnel无法将 UPDATE_BEFORE 和 UPDATE_AFTER 组合成单个 UPDATE 消息,编码时只能将一对更新前/更新后行分别编码为DELETE与INSERT两条 Ogg JSON 消息。也就是说,一个 UPDATE 变更事件经 SeaTunnel 序列化后,在 Kafka 中表现为一条DELETE消息加一条INSERT消息,下游消费者需要自行理解这种"先删后插"的语义(测试断言中可以看到同一id依次输出DELETE与INSERT两条消息)。
库表过滤与容错配置实战
针对多库多表接入的场景,ogg_json.database.include与ogg_json.table.include提供按库、按表的精细订阅能力。例如在 Kafka Source 中同时配置:
source { Kafka { bootstrap.servers = "127.0.0.1:9092" topic = "ogg-all" format = ogg_json ogg_json.database.include = "^PROD" ogg_json.table.include = "^(PRODUCTS|ORDERS)$" schema = { fields { ... } } } }这样只会处理table元字段中库名以PROD开头、表名为PRODUCTS或ORDERS的变更行,其余消息被过滤丢弃,可有效减少下游写入压力。该过滤发生在解析行数据之前,属于 Source 侧的内置能力,无需额外的 Transform 插件。
容错方面,ogg_json.ignore-parse-errors = true可以在个别消息格式异常(如字段缺失、类型不匹配)时跳过出错的行而不是让整个作业失败;出错字段会被置为null。需要权衡的是:开启后数据质量无法保证,适合对完整性不敏感的场景;默认false则保证"坏消息不混入",适合需要严格数据一致性的下游。
端到端验证:E2E 测试配置参考
仓库的 Kafka e2e 测试中提供了两条可直接参考的 Ogg 格式作业配置:
- kafka_source_ogg_to_pgsql.conf:从 Kafka topic
test-ogg-source以format = ogg_json消费,写入 PostgreSQL 目标表(generate_sink_sql = true自动生成写入 SQL,以id为主键); - kafka_source_ogg_to_kafka.conf:Kafka Source 以
ogg_json解析后,由 Kafka Sink 以format = ogg_json重新编码输出到另一个 topictest-ogg-sink,完整演示了"消费 Ogg 消息 → 内部变更行 → 再编码为 Ogg JSON"的转发链路。
这两条配置验证了ogg_json格式同时可用于 Kafka 的 Source 与 Sink 两侧,并且与 JDBC / PostgreSQL、Kafka 等下游存储可以无缝衔接。
注意事项与使用建议
- 字段类型一致性:
schema.fields必须与 Ogg 源表的列定义一致,否则before/after中的值无法正确转换为 SeaTunnel 行类型; - UPDATE / DELETE 依赖变更前镜像:Oracle / PostgreSQL 等数据库需要保证捕获端能够提供
before数据(如 Postgres 需设置REPLICA IDENTITY FULL),否则反序列化会抛出REPLICA IDENTITY相关异常; - UPDATE 的拆分语义:反序列化时一个 UPDATE 事件会展开为两行(
UPDATE_BEFORE+UPDATE_AFTER),序列化时又折叠为两条消息(DELETE+INSERT),下游需针对这种语义设计幂等写入或按主键更新的逻辑; - tombstone 消息:Kafka 中的墓碑消息(空消息体)会被自动跳过,不会产生脏数据;
- 格式归属模块:Ogg JSON 格式实现在 seatunnel-formats/seatunnel-format-json 模块的
ogg子包中,与 Canal JSON、Debezium JSON、Maxwell JSON 等并列,说明该格式适用于 Kafka 等支持format选项的消息类连接器; - 运行前提:使用本格式需要作业能够加载
seatunnel-format-json相关依赖,配置方式与文档 Kafka Source 一致,直接声明format = ogg_json即可。
- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
相关推荐
SeaTunnel ogg_json 格式解析:基于 Oracle GoldenGate 变更日志实现实时数据同步
SeaTunnel ogg_json 格式解析:基于 Oracle GoldenGate 变更日志实现实时数据同步 本篇技术指南围绕 SeaTunnel 的 o
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Ogg Format 实战指南:Oracle GoldenGate JSON Changelog 的读写与解码原理
SeaTunnel Ogg Format 实战指南:Oracle GoldenGate JSON Changelog 的读写与解码原理 本篇技术指南围绕 Sea
数据工程大数据批处理流处理Flink Ogg Format 深度指南:Oracle GoldenGate 变更日志的实时接入与输出
Flink Ogg Format 深度指南:Oracle GoldenGate 变更日志的实时接入与输出 Oracle GoldenGate(简称 Ogg)是
后端大数据流处理批处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考