请问下,flink sql 创建后,源库删除,目标不删除,这个操作有好的解决方法没呀?

请问下,flink sql 创建后,源库删除,目标不删除,这个操作有好的解决方法没呀?

展开
收起
雪哥哥 2022-11-13 20:34:25 3275 分享 版权
4 条回答
写回答
取消 提交回答
  • 一心做个工程师,关注开源和云计算。

    根因:append 流天生没有回撤语义

    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。缺一个,删除都会丢。

    常见漏配点

    • MySQL 侧:binlog_format=ROW、binlog_row_image=FULL 是硬要求。若配成 MINIMAL,-U 事件只带主键和变更列,下游拼不出完整行;-D 事件则只有主键,主键声明不对就匹配不上。
    • Source 侧:未开增量快照(scan.incremental.snapshot.enabled=true),或 server-id、账号权限、snapshot.mode 配错,导致启动阶段就只有 insert。
    • Sink 侧:表没声明主键,或写成 PRIMARY KEY 而非 PRIMARY KEY (...) NOT ENFORCED,Flink 不会切到 upsert 模式,回撤被降级成 append。
    • 有的连接器默认「忽略删除」,参数没打开,-D 事件被静默丢掉。

    方案一:source 用 mysql-cdc,sink 走 upsert

    CREATE TABLE orders_src (
      id BIGINT, amount DECIMAL(10,2),
      PRIMARY KEY (id) NOT ENFORCED
    ) WITH (
      'connector' = 'mysql-cdc',
      'scan.incremental.snapshot.enabled' = 'true',
      ...
    );
    

    关键点三条:

    1. source 表必须声明主键,CDC 才能输出带 op 字段的 changelog。
    2. sink 表主键要和 source 对齐,Flink 由 changelog 模式(UpsertMaterialize)自动下发更新/删除。
    3. Sink 参数按连接器语义开删除能力:JDBC 用 upsert 语义,Doris/StarRocks 关注 sink.enable-delete,Hudi 类湖表要开 changelog.enabled 并把 write.operation 设为 upsert。别照抄参数名,先确认版本和文档语义。

    方案二:整库同步优先用 Flink CDC Pipeline

    如果目标是整库/多表同步,直接用 CDC Pipeline 的 YAML 作业。它在 source、route、sink 三段之间保留 schema 变更和 DML 变更,op=d 的删除事件走的是同一条链路,不需要你手工去对齐每张表的主键和 upsert 配置,落地成本最低。

    兜底:下游确实不允许物理删除

    审计、对账这类场景,目标库本来就不该有物理 DELETE。这时不要硬让它删:直接收回目标账号的 DELETE 权限,表上加强制的 is_deleted / deleted_at 字段,同步逻辑改成把 -D 事件改写为 update 标记。物理删除被阻断在权限层,业务侧查询统一加 is_deleted = 0 过滤。这样数据可追溯,也不会出现「同步删了但没人知道谁删的」。

    排查:先分段,别一上来改 sink

    最省事的定位方式是在 source 后面临时挂一个 print connector:

    • 变更流里能看到 -U / -D,说明 source 和 binlog 没问题,问题在 sink 的表定义或连接器参数。
    • 只有 +I,说明 CDC 根本没吐回撤事件,回去查 binlog_row_image、增量快照开关、账号权限和主键声明。

    多数「删不掉」都不是 Flink 的锅,而是主键没声明或 binlog_row_image 没配成 FULL。

    2026-09-29 23:47:02
    赞同 3 展开评论
  • GitHub https://github.com/co63oc/cloud

    目标库取消用户删除库权限

    2022-11-24 16:10:22
    赞同 28 展开评论
  • 网站:http://ixiancheng.cn/ 微信订阅号:小马哥学JAVA

    与所有 SQL 引擎一样,Flink 查询操作是在表上进行。与传统数据库不同,Flink 不在本地管理静态数据;相反,它的查询在外部表上连续运行。

    Flink 数据处理流水线开始于 source 表。source 表产生在查询执行期间可以被操作的行;它们是查询时 FROM 子句中引用的表。这些表可能是 Kafka 的 topics,数据库,文件系统,或者任何其它 Flink 知道如何消费的系统。 具体可以参考 apache的官方网站

    2022-11-23 09:06:46
    赞同 31 展开评论
  • 十年摸盘键,代码未曾试。 今日码示君,谁有上云事。

    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 表中,而不是将结果加载到系统内存中。

    2022-11-22 13:52:21
    赞同 35 展开评论

实时计算Flink版是阿里云提供的全托管Serverless Flink云服务,基于 Apache Flink 构建的企业级、高性能实时大数据处理系统。提供全托管版 Flink 集群和引擎,提高作业开发运维效率。

收录在圈子:
实时计算 Flink 版(Alibaba Cloud Realtime Compute for Apache Flink,Powered by Ververica)是阿里云基于 Apache Flink 构建的企业级、高性能实时大数据处理系统,由 Apache Flink 创始团队官方出品,拥有全球统一商业化品牌,完全兼容开源 Flink API,提供丰富的企业级增值功能。
还有其他疑问?
咨询AI助理