请问下,flink sql 创建后,源库删除,目标不删除,这个操作有好的解决方法没呀?
版权声明:本文内容由阿里云实名注册用户自发贡献,版权归原作者所有,阿里云开发者社区不拥有其著作权,亦不承担相应法律责任。具体规则请查看《阿里云开发者社区用户服务协议》和《阿里云开发者社区知识产权保护指引》。如果您发现本社区中有涉嫌抄袭的内容,填写侵权投诉表单进行举报,一经查实,本社区将立刻删除涉嫌侵权内容。
INSERT INTO sink SELECT ... FROM src 这类作业,默认把源表当成 append-only 流。源端每来一条数据,Flink 就往 sink 追加一条,整个 DAG 里根本不存在「这条记录对应的是哪条旧数据」的信息。MySQL 里 DELETE 掉一行,binlog 里确实有事件,但普通 JDBC source 只按 SELECT 轮询或全量拉取,看到的是「当前快照」,不产出 -D / -U 变更事件,sink 自然无从删除。
要删得下去,必须两个条件同时成立:source 是 changelog 流,sink 能按主键做 upsert/delete。缺一个,删除都会丢。
binlog_format=ROW、binlog_row_image=FULL 是硬要求。若配成 MINIMAL,-U 事件只带主键和变更列,下游拼不出完整行;-D 事件则只有主键,主键声明不对就匹配不上。scan.incremental.snapshot.enabled=true),或 server-id、账号权限、snapshot.mode 配错,导致启动阶段就只有 insert。PRIMARY KEY 而非 PRIMARY KEY (...) NOT ENFORCED,Flink 不会切到 upsert 模式,回撤被降级成 append。-D 事件被静默丢掉。CREATE TABLE orders_src (
id BIGINT, amount DECIMAL(10,2),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'scan.incremental.snapshot.enabled' = 'true',
...
);
关键点三条:
op 字段的 changelog。UpsertMaterialize)自动下发更新/删除。sink.enable-delete,Hudi 类湖表要开 changelog.enabled 并把 write.operation 设为 upsert。别照抄参数名,先确认版本和文档语义。如果目标是整库/多表同步,直接用 CDC Pipeline 的 YAML 作业。它在 source、route、sink 三段之间保留 schema 变更和 DML 变更,op=d 的删除事件走的是同一条链路,不需要你手工去对齐每张表的主键和 upsert 配置,落地成本最低。
审计、对账这类场景,目标库本来就不该有物理 DELETE。这时不要硬让它删:直接收回目标账号的 DELETE 权限,表上加强制的 is_deleted / deleted_at 字段,同步逻辑改成把 -D 事件改写为 update 标记。物理删除被阻断在权限层,业务侧查询统一加 is_deleted = 0 过滤。这样数据可追溯,也不会出现「同步删了但没人知道谁删的」。
最省事的定位方式是在 source 后面临时挂一个 print connector:
-U / -D,说明 source 和 binlog 没问题,问题在 sink 的表定义或连接器参数。+I,说明 CDC 根本没吐回撤事件,回去查 binlog_row_image、增量快照开关、账号权限和主键声明。多数「删不掉」都不是 Flink 的锅,而是主键没声明或 binlog_row_image 没配成 FULL。
与所有 SQL 引擎一样,Flink 查询操作是在表上进行。与传统数据库不同,Flink 不在本地管理静态数据;相反,它的查询在外部表上连续运行。
Flink 数据处理流水线开始于 source 表。source 表产生在查询执行期间可以被操作的行;它们是查询时 FROM 子句中引用的表。这些表可能是 Kafka 的 topics,数据库,文件系统,或者任何其它 Flink 知道如何消费的系统。 具体可以参考 apache的官方网站
SQL客户端是一个交互式的客户端,用于向Flink提交SQL查询并将结果可视化。
与所有SQL引擎一样,Flink查询操作是在表上进行。 与传统数据库不同,Flink不在本地管理静态数据;相反,它的查询在外部表上连续运行。
可以通过SQL客户端或使用环境配置文件来定义表。 SQL客户端支持类似于传统的SQL DDL命令, 用于创建,修改,删除表。 Flink支持不同的连接器和格式相结合以定义表。
Flink SQL与传统数据库查询的不同之处在于,Flink SQL持续消费到达的行并对其结果进行更新。 一个连续查询永远不会终止,并会产生一个动态表作为结果。 动态表是Flink中Table API和SQL对流数据支持的核心概念。
Sink表 当运行此查询时,SQL客户端实时但是以只读方式提供输出。 存储结果,作为报表或仪表板的数据来源,需要写到另一个表。 这可以使用INSERT INTO语句来实现。 INSERT INTO语句将作为一个独立查询被提交到Flink集群中。
提交后,它将运行并将结果直接存储到sink 表中,而不是将结果加载到系统内存中。
实时计算Flink版是阿里云提供的全托管Serverless Flink云服务,基于 Apache Flink 构建的企业级、高性能实时大数据处理系统。提供全托管版 Flink 集群和引擎,提高作业开发运维效率。