Flink - Kafka Connector EXACTLY ONCE error

Viewed 108

I'm using Flink 1.15.0. For Kafka integration I wrote

KafkaSource:

public static KafkaSource<String> kafkaSource(String bootstrapServers, String topic, String groupId) {
    return KafkaSource.<String>builder()
            .setBootstrapServers(bootstrapServers)
            .setTopics(topic)              
            .setGroupId(groupId)                
            .setStartingOffsets(OffsetsInitializer.latest())
            .setValueOnlyDeserializer(new SimpleStringSchema())
            .build();

and KafkaSink:

public static KafkaSink<WordCountPojo> kafkaSink(String brokers, String topic, Properties producerProperties) {
    return KafkaSink.<WordCountPojo>builder()
            .setKafkaProducerConfig(producerProperties)
            .setBootstrapServers(brokers)
            .setRecordSerializer(new WordCountPojoSerializer(topic))
            .setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
            .setTransactionalIdPrefix(UUID.randomUUID().toString())
            .build();
}

for completeness this is my custom serializer

public class WordCountPojoSerializer implements KafkaRecordSerializationSchema<WordCountPojo> {
    private String topic;
    private ObjectMapper mapper;

    public WordCountPojoSerializer(String topic) {
        this.topic = topic;
        mapper = new ObjectMapper();
    }

    @Override
    public ProducerRecord<byte[], byte[]> serialize(WordCountPojo wordCountPojo, KafkaSinkContext kafkaSinkContext, Long timestamp) {
        try {
            byte[] serializedValue = mapper.writeValueAsBytes(wordCountPojo);
            return new ProducerRecord<>(topic, null,  timestamp,null, serializedValue);
        } catch (JsonProcessingException e) {
            return null;
        }
    }
}

Property transaction.timeout.ms is set to 60_000ms (1 minute)

My application run for a while, then, suddenly, stop with the exception

Caused by: org.apache.kafka.common.errors.InvalidProducerEpochException: Producer attempted to produce with an old epoch.

It seems that producer tries to produce record between two transaction, but I'm not sure. I could increase transaction.timeout.ms to max 900_000 (15 minute) but I don't think this may solve the problem.

Can you help me please?

0 Answers
Related