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?