Apache Flink Job execution fails whenever delete change log in received from Debezium MySql CDC

Viewed 284

I have implemented Debezium MySql connector with Apache Flink. I am able to receive CDC event through Kafka. Insert and update event are working fine.But as soon as delete event is received, Flink job and application stops working with Job execution failed exception error. Delete event is received in my code and sink invoke method is getting triggered and my delete query is running successfully but after completion Flink job is stopped with error message Job execution failed. Not aware if it is a bug in Flink or any configuration is not right. I have attached my Flink main configuration class and Sink class.

**//Apache Flink job configuration and implementation class**



final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    
            try{
    
                Properties props = new Properties();
                props.put("bootstrap.servers", "172.17.0.3:9092");
                props.put("zookeeper.connect", "172.17.0.2:2181");
                props.put("group.id", "metric-group");
    //            props.put("value.serializer", AppConstant.VALUE_SERIALIZER);
                props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
                props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
                props.put("auto.offset.reset", "latest");
    
                DataStreamSource<String> dataStreamSource = env.addSource(new FlinkKafkaConsumer<>(
                        "dbserver1.inventory.customers",
                        new SimpleStringSchema(),
                        props));
    
    
                DataStream<MySqlCDCResponse> testEntity = dataStreamSource
                        .map(string -> new Gson().fromJson(string, MySqlCDCResponse.class));
    
    
                testEntity.addSink(new MySqlSink());
    
                System.out.println("Data::");
    
                env.execute("MySqlApp");
    
            }
            catch (Exception ex){
                System.out.println("Error -> " + ex.getMessage());
            }
    
    
**//MySql Sink class that I am using with Apache Flink **


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 ps;
    private Connection connection;
//    private static final Logger log = LoggerFactory.getLogger(MySqlSink.class);


    private static String DELETE_OPERATION = "select * from customers where email=?";

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        connection = getConnection();

        String sql = "insert into customers (first_name, last_name, email) values( ?, ?, ?) " +
                "ON DUPLICATE KEY UPDATE first_name = VALUES(first_name) , last_name = VALUES(last_name);";
        if (connection != null) {
            ps = this.connection.prepareStatement(sql);
        }
    }

    @Override
    public void close() throws Exception {
        super.close();

        if (connection != null) {
            connection.close();
        }
        if (ps != null) {
            ps.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 (ps == null) {
            return;
        }
        try {
            if(Objects.nonNull(value.getPayload().getBefore()) && Objects.isNull(value.getPayload().getAfter())) {
                //connection = getConnection();
                System.out.println("Delete Operation " + value.getPayload().toString());
                ps = this.connection.prepareStatement(DELETE_OPERATION);
                ps.setString(1, value.getPayload().getBefore().getEmail());
                ps.executeQuery();
            }
            else {
                ps.setString(1, value.getPayload().getAfter().getFirst_name());
                ps.setString(2, value.getPayload().getAfter().getLast_name());
                ps.setString(3, value.getPayload().getAfter().getEmail());
                ps.executeUpdate();
            }
        }
        catch (Exception ex) {
            System.out.println("Error:: " + 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) {
//            log.error("-----------mysql get connection has exception , msg = {}", e.getMessage());
            System.out.println("Connection failure::" + ex.getMessage());
        }
        return con;
    }

}

public class MySqlCDCResponse {

private int id;
private String first_name;
private String last_name;
private String email;

public MySqlCDCResponse(String first_name,  String last_name, String email){
    this.first_name =first_name;
    this.last_name =last_name;
    this.email =email;
}

private PayloadBean payload;

public PayloadBean getPayload() {
    return payload;
}

public void setPayload(PayloadBean payload) {
    this.payload = payload;
}

public static class PayloadBean {

    private BeforeBean before;
    private AfterBean after;
    private SourceBean source;
    private String op;
    private long ts_ms;
    private Object transaction;

    public BeforeBean getBefore() {
        return before;
    }

    public void setBefore(BeforeBean before) {
        this.before = before;
    }

    public AfterBean getAfter() {
        return after;
    }

    public void setAfter(AfterBean after) {
        this.after = after;
    }

    public SourceBean getSource() {
        return source;
    }

    public void setSource(SourceBean source) {
        this.source = source;
    }

    public String getOp() {
        return op;
    }

    public void setOp(String op) {
        this.op = op;
    }

    public long getTs_ms() {
        return ts_ms;
    }

    public void setTs_ms(long ts_ms) {
        this.ts_ms = ts_ms;
    }

    public Object getTransaction() {
        return transaction;
    }

    public void setTransaction(Object transaction) {
        this.transaction = transaction;
    }

    public static class BeforeBean {

        private int id;
        private String first_name;
        private String last_name;
        private String email;

        public int getId() {
            return id;
        }

        public void setId(int id) {
            this.id = id;
        }

        public String getFirst_name() {
            return first_name;
        }

        public void setFirst_name(String first_name) {
            this.first_name = first_name;
        }

        public String getLast_name() {
            return last_name;
        }

        public void setLast_name(String last_name) {
            this.last_name = last_name;
        }

        public String getEmail() {
            return email;
        }

        public void setEmail(String email) {
            this.email = email;
        }

        @Override
        public String toString() {
            return "BeforeBean{" +
                    "id=" + id +
                    ", first_name='" + first_name + '\'' +
                    ", last_name='" + last_name + '\'' +
                    ", email='" + email + '\'' +
                    '}';
        }
    }

    public static class AfterBean {

        private int id;
        private String first_name;
        private String last_name;
        private String email;

        public int getId() {
            return id;
        }

        public void setId(int id) {
            this.id = id;
        }

        public String getFirst_name() {
            return first_name;
        }

        public void setFirst_name(String first_name) {
            this.first_name = first_name;
        }

        public String getLast_name() {
            return last_name;
        }

        public void setLast_name(String last_name) {
            this.last_name = last_name;
        }

        public String getEmail() {
            return email;
        }

        public void setEmail(String email) {
            this.email = email;
        }

        @Override
        public String toString() {
            return "AfterBean{" +
                    "id=" + id +
                    ", first_name='" + first_name + '\'' +
                    ", last_name='" + last_name + '\'' +
                    ", email='" + email + '\'' +
                    '}';
        }
    }

    public static class SourceBean {

        private String version;
        private String connector;
        private String name;
        private long ts_ms;
        private String snapshot;
        private String db;
        private Object sequence;
        private String table;
        private int server_id;
        private Object gtid;
        private String file;
        private int pos;
        private int row;
        private Object thread;
        private Object query;

        public String getVersion() {
            return version;
        }

        public void setVersion(String version) {
            this.version = version;
        }

        public String getConnector() {
            return connector;
        }

        public void setConnector(String connector) {
            this.connector = connector;
        }

        public String getName() {
            return name;
        }

        public void setName(String name) {
            this.name = name;
        }

        public long getTs_ms() {
            return ts_ms;
        }

        public void setTs_ms(long ts_ms) {
            this.ts_ms = ts_ms;
        }

        public String getSnapshot() {
            return snapshot;
        }

        public void setSnapshot(String snapshot) {
            this.snapshot = snapshot;
        }

        public String getDb() {
            return db;
        }

        public void setDb(String db) {
            this.db = db;
        }

        public Object getSequence() {
            return sequence;
        }

        public void setSequence(Object sequence) {
            this.sequence = sequence;
        }

        public String getTable() {
            return table;
        }

        public void setTable(String table) {
            this.table = table;
        }

        public int getServer_id() {
            return server_id;
        }

        public void setServer_id(int server_id) {
            this.server_id = server_id;
        }

        public Object getGtid() {
            return gtid;
        }

        public void setGtid(Object gtid) {
            this.gtid = gtid;
        }

        public String getFile() {
            return file;
        }

        public void setFile(String file) {
            this.file = file;
        }

        public int getPos() {
            return pos;
        }

        public void setPos(int pos) {
            this.pos = pos;
        }

        public int getRow() {
            return row;
        }

        public void setRow(int row) {
            this.row = row;
        }

        public Object getThread() {
            return thread;
        }

        public void setThread(Object thread) {
            this.thread = thread;
        }

        public Object getQuery() {
            return query;
        }

        public void setQuery(Object query) {
            this.query = query;
        }

        @Override
        public String toString() {
            return "SourceBean{" +
                    "version='" + version + '\'' +
                    ", connector='" + connector + '\'' +
                    ", name='" + name + '\'' +
                    ", ts_ms=" + ts_ms +
                    ", snapshot='" + snapshot + '\'' +
                    ", db='" + db + '\'' +
                    ", sequence=" + sequence +
                    ", table='" + table + '\'' +
                    ", server_id=" + server_id +
                    ", gtid=" + gtid +
                    ", file='" + file + '\'' +
                    ", pos=" + pos +
                    ", row=" + row +
                    ", thread=" + thread +
                    ", query=" + query +
                    '}';
        }
    }

    @Override
    public String toString() {
        return "PayloadBean{" +
                "before=" + before +
                ", after=" + after +
                ", source=" + source +
                ", op='" + op + '\'' +
                ", ts_ms=" + ts_ms +
                ", transaction=" + transaction +
                '}';
    }
}

@Override
public String toString() {
    return "MySqlCDCResponse{" +
            "payload=" + payload +
            '}';
}

}

0 Answers
Related