Flink Ogg Format 实战:基于 Oracle GoldenGate JSON 的 Changelog 数据接入指南
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载OggOracle GoldenGateFormat 是 Apache Flink 提供的一种 Changelog-Data-CaptureCDC格式它允许 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 并把变更事件解释为 changelogCREATE 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 RowKindIinsertINSERT取after字段UupdateUPDATE_BEFOREbefore字段UPDATE_AFTERafter字段两条消息DdeleteDELETE取before字段Ttruncate其他未知值时若未开启忽略解析错误则抛出异常对于Uupdate与Ddelete操作如果before字段为 null反序列化会抛出IllegalStateException。源码中给出的排查提示是如果使用 Ogg Postgres Connector需要确认 Postgres 表已设置REPLICA IDENTITY为FULL级别否则无法拿到变更前的镜像数据OggJsonDeserializationSchema。另外反序列化遇到 null 或空字节数组Kafka 的 tombstone 消息时会直接跳过不会产出任何记录。可用元数据Available Metadataogg-json格式可以将以下格式元数据暴露为表定义中的只读VIRTUAL列Key数据类型描述tableSTRING NULL完全限定的表名格式为目录名.模式名.表名CATALOG NAME.SCHEMA NAME.TABLE NAMEprimary-keysARRAYSTRING 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-ddTHH: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 ARRAYSTRING 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 Optionsogg-json格式支持以下配置选项均定义于 OggJsonFormatFactory 与 OggJsonFormatOptions 中选项是否必填默认值类型描述format必填无String指定使用的格式此处应为ogg-jsonogg-json.ignore-parse-errors可选falseBoolean遇到解析错误时跳过对应字段和行而不是失败出错时字段会被置为 nullogg-json.timestamp-format.standard可选SQLString指定输入/输出的时间戳格式目前支持SQL与ISO-8601两种取值详见下方说明ogg-json.map-null-key.mode可选FAILString序列化 Map 数据时遇到 null key 的处理模式支持FAIL、DROP、LITERALogg-json.map-null-key.literal可选nullString当ogg-json.map-null-key.mode为LITERAL时用于替换 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仅支持SQL和ISO-8601ogg-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 到 RowDataogg-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 JSONogg-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_AFTERbefore置为 nullafter写入当前行数据op_type置为IUPDATE_BEFORE/DELETEbefore写入当前行数据after置为 nullop_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 会以先 DELETEbefore后 INSERTafter两条 Ogg 消息落盘下游消费者需要按此语义还原变更。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/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 是阿里巴巴开源的 CDCC大数据流处理批处理数据工程Flink 集成 Maxwell JSON 格式基于 MySQL CDC 的 Changelog 流接入与输出完整指南Flink 集成 Maxwell JSON 格式基于 MySQL CDC 的 Changelog 流接入与输出完整指南 Maxwell 是业界常用的 CDC大数据流处理批处理数据工程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻

漫威电影观影顺序2026最新避坑指南:别再用时间线坑自己了

漫威电影观影顺序2026最新避坑指南:别再用时间线坑自己了

漫威电影观影顺序2026最新避坑指南:别再用时间线坑自己了 面试被问原理答不上来,这种尴尬感就像你拿着《复仇者联盟4》的截图去跟HR聊剧情,对方一脸懵逼。很多转行做内容策划、影视数据分析或后端开发的伙伴,都栽在“观影顺序”这个看似简单实则充…

2026/9/23 14:36:41 阅读更多 →
Triton Inference Server 测试指南:QA 模型仓库生成、镜像构建与 L0 测试全流程

Triton Inference Server 测试指南:QA 模型仓库生成、镜像构建与 L0 测试全流程

Triton Inference Server 测试指南:QA 模型仓库生成、镜像构建与 L0 测试全流程 【免费下载链接】server The Triton Inference Server provides an optimized cloud and edge inferencing solution. 项目地址: https://gitcode.com/gh_mirrors/server117/server…

2026/9/23 14:35:40 阅读更多 →
SIwave SI/PI协同仿真指南:从IR Drop到S参数的高速PCB仿真

SIwave SI/PI协同仿真指南:从IR Drop到S参数的高速PCB仿真

简介:面向信号完整性与电源完整性工程师的SIwave中文培训手册,围绕高速PCB设计中SI/PI与EMI的关联展开,既讲清传输线、特性阻抗、S参数、去耦电容、目标阻抗等基础理论,又结合完整仿真实例演示如何排查谐振、阻抗失配、SSN、DC压降…

2026/9/23 14:35:40 阅读更多 →

最新新闻

信号分析与处理实验全链路:从采样到滤波器设计的MATLAB实现

信号分析与处理实验全链路:从采样到滤波器设计的MATLAB实现

简介:这份资源是南京邮电大学「信号分析与处理实验」课程的完整实验报告,面向正在修读数字信号处理、信号与系统相关课程的高校学生,以及需要借助 MATLAB 完成实验与课程设计的自学者。报告覆盖信号的产生和运算、连续时间信号的频域分析、信…

2026/9/23 23:38:58 阅读更多 →
三款智能颈椎与腰部牵引理疗仪硬件横评:仿生揉捏与气压热敷实测

三款智能颈椎与腰部牵引理疗仪硬件横评:仿生揉捏与气压热敷实测

三款智能颈椎与腰部牵引理疗仪硬件横评:仿生揉捏与气压热敷实测秋分过后气温骤降,长期坐在电脑前写代码的开发者与上了年纪的长辈,最容易遭遇颈椎僵硬、肩背酸痛与腰椎间盘劳损的集中爆发: 老爸年轻时当老师落下了颈椎病&#xff…

2026/9/23 23:38:58 阅读更多 →
南山一经深度拆解:从异兽到祭祀,读懂山海经的博物志密码

南山一经深度拆解:从异兽到祭祀,读懂山海经的博物志密码

1. 为什么我要逐字啃完南山一经《山海经》第一卷南山经里的南山一经,全文不过几百字,却藏着四十多座山、十几种异兽、一堆矿产和祭祀规矩。很多人翻《山海经》都是跳着看,专挑九尾狐、凤凰这些网红神兽,但真正想把这本书读透的人&…

2026/9/23 23:38:58 阅读更多 →
大模型并不是真正的记忆:从神经元突触重塑看权重的冷热之分

大模型并不是真正的记忆:从神经元突触重塑看权重的冷热之分

大模型并不是真正的记忆:从神经元突触重塑看权重的冷热之分昨天老妈在厨房里找东西时,发生了一幕全家人都极其熟悉的生活小插曲:老妈站在调料架前,拍了拍脑门:"哎呀!我刚才明明记得把新买的白胡椒粉放…

2026/9/23 23:38:58 阅读更多 →
长辈友好型节气动态插画工程:纯 SVG 矢量绘制与轻量 CSS 路径动画

长辈友好型节气动态插画工程:纯 SVG 矢量绘制与轻量 CSS 路径动画

长辈友好型节气动态插画工程:纯 SVG 矢量绘制与轻量 CSS 路径动画在很多针对长辈的节气提醒与家庭生活看板中,工程师为了展示节气氛围,常常直接在页面中嵌入体积庞大的 GIF 动图或 MP4 短视频。 但在家庭低功耗平板、电子相框或老式电视盒子上…

2026/9/23 23:38:58 阅读更多 →
OK交易所Python API封装实战:现货、杠杆与历史数据调用指南

OK交易所Python API封装实战:现货、杠杆与历史数据调用指南

简介:这份Python资源包围绕OKEx交易所Web API的调用展开,面向希望用代码接入加密货币市场的开发者与量化交易初学者。内容覆盖杠杆交易、现货交易、历史记录与历史数据获取等核心场景,并涉及MVC架构下的应用组织方式,适合具备Pyth…

2026/9/23 23:37:58 阅读更多 →

日新闻

3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A…

2026/9/23 0:00:23 阅读更多 →
2k显示屏性能优化踩坑:版本升级后API全变了,这份源码解析救了我

2k显示屏性能优化踩坑:版本升级后API全变了,这份源码解析救了我

2k显示屏性能优化踩坑:版本升级后API全变了,这份源码解析救了我 刚把开发环境的显示器从1080P换到2K,跑老项目直接报错,版本升级后 API…

2026/9/23 0:01:25 阅读更多 →
3步搞定美眉图实战项目,告别官方文档抓不住重点

3步搞定美眉图实战项目,告别官方文档抓不住重点

3步搞定美眉图实战项目,告别官方文档抓不住重点 官方文档翻了三遍还是云里雾里?别急,美眉图在实战项目中常被用来做数据可视化,但它的原理比你想的简单。今天咱们直接上手,用一个完整的小项目把美眉图跑通,不再死磕那些冗长的理论说明。…

2026/9/23 0:01:25 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/23 4:55:02 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/23 4:49:06 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/23 9:53:41 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/23 9:53:40 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/23 9:53:40 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/23 9:53:40 阅读更多 →