接收到Debezium MySQL CDC删除事件时Apache Flink作业执行失败如何解决
Flink CDC删除事件导致作业崩溃问题修复方案
核心错误点
- 预编译对象被非法修改:你在
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
相关产品推荐
相关产品推荐

