You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

接收到Debezium MySQL CDC删除事件时Apache Flink作业执行失败如何解决

核心错误点

  • 预编译对象被非法修改:你在open()阶段初始化的全局PreparedStatement是插入/更新语句的预编译实例,处理删除事件时直接对该全局变量重新赋值为查询语句实例,下一条插入/更新事件到来时,会用查询语句的预编译对象设置插入参数、执行executeUpdate(),直接触发SQL异常导致作业失败
  • 删除逻辑错误:你定义的删除操作SQL是查询语句,根本不会执行数据删除,且查询产生的结果集未关闭,会导致JDBC资源泄漏

修复后的MySqlSink代码

package entity;

import dto.MySqlCDCResponse;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.util.Objects;

public class MySqlSink extends RichSinkFunction<MySqlCDCResponse> {
    // 全局只保留插入更新的预编译对象,运行过程中不要修改
    private PreparedStatement upsertPs;
    private Connection connection;

    private static String UPSERT_SQL = "insert into customers (first_name, last_name, email) values( ?, ?, ?) " +
            "ON DUPLICATE KEY UPDATE first_name = VALUES(first_name) , last_name = VALUES(last_name);";
    // 替换为真正的删除SQL
    private static String DELETE_SQL = "delete from customers where email=?";

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        connection = getConnection();
        if (connection != null) {
            upsertPs = this.connection.prepareStatement(UPSERT_SQL);
        }
    }

    @Override
    public void close() throws Exception {
        super.close();
        if (upsertPs != null) {
            upsertPs.close();
        }
        if (connection != null) {
            connection.close();
        }
    }

    @Override
    public void invoke(MySqlCDCResponse value, Context context) throws Exception {
        System.out.println("Thread name:: " + Thread.currentThread().getName());
        System.out.println("value= " + value.getPayload().toString());
        if (connection == null) {
            return;
        }
        try {
            if(Objects.nonNull(value.getPayload().getBefore()) && Objects.isNull(value.getPayload().getAfter())) {
                System.out.println("Delete Operation " + value.getPayload().toString());
                // 删除操作单独创建预编译对象,用完自动关闭
                try(PreparedStatement deletePs = connection.prepareStatement(DELETE_SQL)) {
                    deletePs.setString(1, value.getPayload().getBefore().getEmail());
                    deletePs.executeUpdate();
                }
            } else {
                if (upsertPs == null) {
                    return;
                }
                upsertPs.setString(1, value.getPayload().getAfter().getFirst_name());
                upsertPs.setString(2, value.getPayload().getAfter().getLast_name());
                upsertPs.setString(3, value.getPayload().getAfter().getEmail());
                upsertPs.executeUpdate();
            }
        } catch (Exception ex) {
            System.out.println("处理数据异常:: " + ex.getMessage());
        }
    }

    private static Connection getConnection() {
        Connection con = null;
        try {
            Class.forName("com.mysql.jdbc.Driver");
            con = DriverManager.getConnection("jdbc:mysql://localhost:3306/test?useUnicode=true&characterEncoding=UTF-8", "root", "root@123");
        } catch (Exception ex) {
            System.out.println("数据库连接失败::" + ex.getMessage());
        }
        return con;
    }
}

额外优化建议

  • 建议后续用Debezium CDC事件里的op字段判断操作类型(d代表删除,c代表插入,u代表更新),比判断before/after是否为空更准确
  • 生产环境不要直接用RichSinkFunction写JDBC逻辑,建议用Flink官方的JDBC Connector,自带连接池、重试、Exactly-Once语义保证,稳定性更高

内容的提问来源于stack exchange,提问作者chandra rautela

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.30 20:54:03