ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Flink Ogg Format 实战:基于 Oracle GoldenGate JSON 的 Changelog 数据接入指南

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),仅供参考
返回列表