How to skip kafka history data in flink job if certain lag is encountered?

Viewed 131

Sometimes we encounter lag in kafka consumer due to some external issues.

Flink job will always consume kafka history (delayed data) with exactly-once semantics, but here's a scenario:

We will skip delayed data when kafka consumer lag is too much in order to let our downstream service get the latest data in time. I am thinking to set a window period to do it. What should I code for it?

2 Answers

You could stop the Flink job and use kafka-consumer-groups CLI from Kafka to seek the consumer group forward (assuming Flink is using one, rather than maintaining offsets itself)

When the job restarts, it'll start from the new offset location

I'd say your least painful option is to always read all the messages, but do not process (discard them as soon as possible) the ones you want to skip. Just reading and discarding without any further processing is really fast.

Related