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);