why kafka producer has been closed for each commit in flink EXACTLY_ONCE mode

Viewed 584

i use flink-connector-kafka in my flink application, the semantic is set to EXACTLY_ONCE, and i see the log keep printing the kafka has been closed and reconnect like this :

Closing the Kafka producer with timeoutMillis = 0 ms.
Proceeding to force close the producer since pending requests could not be completed within timeout 0 ms.

and i view the source code, found the close call from producer commit function, the commit fun call the recycleTransactionalProducer in finally block, and the recycleTransactionalProducer fun call the close fun, whitch print the log, so why the kafka producer has been closed for each commit?

the source code from package :

org.apache.flink.streaming.connectors.kafka;

org.apache.kafka.clients.producer;

    @Override
    protected void commit(FlinkKafkaProducer.KafkaTransactionState transaction) {
        if (transaction.isTransactional()) {
            try {
                transaction.producer.commitTransaction();
            } finally {
                recycleTransactionalProducer(transaction.producer);
            }
        }
    }

    private void recycleTransactionalProducer(FlinkKafkaInternalProducer<byte[], byte[]> producer) {
        availableTransactionalIds.add(producer.getTransactionalId());
        producer.flush();
        producer.close(Duration.ofSeconds(0));
    }

    private void close(Duration timeout, boolean swallowException) {
        long timeoutMs = timeout.toMillis();
        if (timeoutMs < 0)
            throw new IllegalArgumentException("The timeout cannot be negative.");
        log.info("Closing the Kafka producer with timeoutMillis = {} ms.", timeoutMs);

1 Answers

Quoting from http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/Problems-with-FlinkKafkaProducer-closing-after-timeoutMillis-9223372036854775807-ms-td39488.html :

... when using exactly-once semantics for the FlinkKafkaProducer, there is a fixed-sized pool of short-living Kafka producers that are created for each concurrent checkpoint. When a checkpoint begins, the FlinkKafkaProducer creates a new producer for that checkpoint. Once said checkpoint completes, the producer for that checkpoint is attempted to be closed and recycled. So, it is normal to see logs of Kafka producers being closed if you're using an exactly-once transactional FlinkKafkaProducer.

Related