
数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本文基于 SeaTunnel 官方中文文档docs/zh/connector-v2/formats/ogg-json.md结合seatunnel-formats/seatunnel-format-json模块中 Ogg JSON 反序列化/序列化源码与端到端测试配置系统讲解 OggOracle GoldenGateJSON 格式在 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_typeIINSERT、UUPDATE、DDELETE决定消息被解释为插入、更新还是删除事件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_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兼容从源码可以确认这些选项的解析方式ogg_json.ignore-parse-errors对应JsonFormatOptions.IGNORE_PARSE_ERRORS默认值为false即默认任何解析错误都会使作业失败ogg_json.database.include与ogg_json.table.include均无默认值noDefaultValue未配置时表示不过滤任何数据库 / 表两个 include 正则最终会被编译为 JavaPatternOggJsonDeserializationSchema.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为intname/description为stringweight为doubleSeaTunnel 据此把before/after中的 JSON 对象转换为行数据job.mode STREAMING配合start_mode earliest表示以流式模式从最早位点开始持续消费反序列化时op_type为I的消息产生INSERTI行为U的消息同时产生UPDATE_BEFORE-U与UPDATE_AFTERU两行为D的消息产生DELETE-D行下游 JDBC Sink 会根据行的 RowKind 执行对应的写入语义。变更消息的反序列化语义Ogg JSON 反序列化器位于 OggJsonDeserializationSchema.java其核心逻辑deserializeMessage方法遵循以下流程跳过墓碑消息如果消息为null或长度为 0Kafka 的 tombstone 记录直接返回不产生任何输出行解析 JSON将字节流解析为 Jackson 的ObjectNode解析失败时若配置了ogg_json.ignore-parse-errors true则静默跳过否则抛出SeaTunnelRuntimeException数据库 / 表过滤若配置了 include 正则按上文所述方式匹配table元字段拆分出的库名与表名不匹配的消息被丢弃按操作类型分发IINSERT读取after节点转换为行以RowKind.INSERT输出UUPDATE读取before与after节点分别以RowKind.UPDATE_BEFORE和RowKind.UPDATE_AFTER输出两行若before为空则抛出IllegalStateExceptionDDELETE读取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 XXtestFilteringTables通过setDatabase(^OG.*).setTable(^TBL.*)验证库表正则过滤完整的序列化/反序列化往返断言UPDATE 事件在反序列化时拆成-U与U两行在序列化时又分别编码为{type:DELETE}与{type:INSERT}。变更消息的序列化将 SeaTunnel 变更流编码为 Ogg JSONSeaTunnel 同样支持把内部的 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输出的 OggtypeINSERTINSERTUPDATE_AFTERINSERTUPDATE_BEFOREDELETEDELETEDELETE其他抛出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 topictest-ogg-source以format ogg_json消费写入 PostgreSQL 目标表generate_sink_sql true自动生成写入 SQL以id为主键kafka_source_ogg_to_kafka.confKafka 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_BEFOREUPDATE_AFTER序列化时又折叠为两条消息DELETEINSERT下游需针对这种语义设计幂等写入或按主键更新的逻辑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即可。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐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是后端大数据流处理批处理上一篇CANN驱动ECC时间查询下一篇nvitop终极指南如何高效监控GPU加速的影视特效渲染任务创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考