news 2026/8/25 23:55:31

FlinkSQL 处理 binlog:changelog 的三种处理方式

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
FlinkSQL 处理 binlog:changelog 的三种处理方式

「我的数据空间」实时计算实践笔记 · Flink SQL 系列

使用 Flink 临时表

使用 DDL 声明 对应的 schema 和 format

CREATETABLEKafkaSource(idVARCHAR,`count`BIGINT,changelogBOOLEAN)with('topic'='topic_ods_order_event','connector'='kafka','format'='binlog','binlog.with-changelog'='true','properties.bootstrap.servers'='kafka-bootstrap:9092','scan.startup.mode'='latest-offset');

说明:

  1. changelog 是一个固定的字段名, 用于表示消息的属性, true代表增加, false代表删除, changelog必须在 ‘connector.with-changelog’ = 'true’时才会生效, 否则,chanlog会被当作一个普通字段, 如果原始的mysql表中不包含这个changelog字段,则会报错.
  2. 如果原始表中已经存在了changelog这个字段,且设置了changelog字段, changelog字段会优先作为消息的属性信息,而不是原始的字段, 为避免冲突,可以设置 ‘connector.changelog-name’ = ‘xxx’ 来修改用于存放changelog的字段名.

除去changelog之外,还支持获得binlog的其他属性

字段名说明配置参数配置字段名
changelogbinlog的消息属性‘connector.with-changelog’‘connector.changelog-name’
offsetbinlog的offset‘connector.with-offset’‘connector.offset-name’
binlogTimebinlog产生的时间‘connector.with-timestamp’‘connector.timestamp-name’

changelog的常用处理方式:

3.1 直接过滤掉删除的记录

SELECT*FROMKafkaSourceWHEREchangelog=true;

这种情况只适合于没有主键删除的情况, 只需要处理add的消息即可.

3.2 将changelog字段作为普通字段处理

SELECTid,LAST_VALUE(`count`),LAST_VALUE(changelog)FROMKafkaSourceGROUPBYid;

在FlinkSQL这一层不处理changelog,而是将changelog当作普通字段来处理, 并写入到下游系统, 由下游的系统来处理.
原始数据:

+0001,1,true+0001,2,true+0001,2,false

经过处理之后

+0001,1,true-0001,2,true+0001,2,true-0001,2,false+0001,2,false

经过以上的处理, 最终 + 0001, 2, false这条记录会被写入到最终的sink表.

典型应用场景:

binlog数据实时导入到iceberg中:

CREATETABLEKafkaSource(idVARCHAR,`count`BIGINT,changelogBOOLEAN)with('topic'='topic_ods_order_event','connector'='kafka','format'='binlog','binlog.with-changelog'='true','properties.bootstrap.servers'='kafka-bootstrap:9092','scan.startup.mode'='latest-offset');insertintoiceberg_catalog.dw.dwd_orderselectevent_guidasguid,event_typeastypefromKafkaSourcegroupbyevent_guid,event_type;

本文收录于「我的数据空间」技术库——一套可私有化部署的数据平台(数据集成 / 实时计算 / 数据湖 / 湖仓查询 / 智能问数)。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/25 23:46:17

AI智能体系统安全防御:从间接提示词注入到架构级防护

1. 从一次真实的“安全策略拦截”事件谈起最近在调试一个多智能体系统时,遇到了一个让我印象深刻的报错。系统日志里赫然写着:ef1 usb device news disk hans been blocked by zhe current security policy。这行看似语法混乱、充满拼写错误的英文&#…

作者头像 李华
网站建设 2026/8/25 23:45:46

Agent找工作有感,说点搜不到的(已入职)

我从投简历到入职,折腾了差不多两个月,现在坐在工位上回头看这段求职经历,有些话真的是网上搜不到的,想写出来给正在找agent开发岗的朋友一点参考。 📌先说一个我踩过最大的坑。 面经都在告诉你agent岗火、薪资高、缺…

作者头像 李华
网站建设 2026/8/25 23:43:12

159、洞察驱动的实战标题——Android Camera HAL3状态机深度解析——从request到result的每一毫秒延迟来源

159、洞察驱动的实战标题——Android Camera HAL3状态机深度解析——从request到result的每一毫秒延迟来源 上周三凌晨两点,客户现场反馈:某旗舰机型在暗光预览下,取景画面出现周期性“卡顿感”,每三秒左右一次,每次持续约两百毫秒。抓了log,发现预览帧率在卡顿瞬间从30…

作者头像 李华
网站建设 2026/8/25 23:40:59

10 极物科技 | DALI调光灯 - DT8 双色温控制专题

极物科技 | DALI调光灯 - DT8 双色温控制专题 前言 极物科技致力于智能控制系统及硬件研发,以"极物OS多协议融合"打破生态壁垒。 一句话概述:本文深入解析 DALI DT8 双色温控制的完整技术实现,帮助技术人员掌握双色温调光的核心协议…

作者头像 李华
网站建设 2026/8/25 23:40:43

LLM、Agent与RAG:从核心原理到工程实践,构建可靠AI应用

1. 从“聊天机器人”到“智能大脑”:AI术语的迷雾与现实如果你最近刷社交媒体、看科技新闻,或者只是和同事朋友聊天,大概率会频繁听到“LLM”、“Agent”、“RAG”这几个词。它们被包装成各种神奇的概念,仿佛有了它们,…

作者头像 李华