I have question regarding handling of consumers death due to exceeding the timeout values.
my example configuration:
session.timeout.ms = 10000 (10 seconds)
heartbeat.interval.ms = 2000 (2 seconds)
max.poll.interval.ms = 300000 (5 minutes)
I have 1 topic, 10 partitions, 1 consumer group, 10 consumers (1 partition = 1 consumer).
From my understanding consuming messages in Kafka, very simplified, works as follows:
- consumer polls 100 records from topic
- a
heartbeatsignal is sent to broker - processing records in progress
- processing records completes
- finalize processing (commit, do nothing etc.)
- repeat #1-5 in a loop
My question is, what happens if time between heartbeats takes longer than previously configured session.timeout.ms. I understand the part, that if session times out, the broker initializes a re-balance, the consumer which processing took longer than the session.timeout.ms value is marked as dead and a different consumer is assigned/subscribed to that partition.
Okey, but what then...?
- Is that long-processing consumer removed/unsubscribed from the topic and my application is left with
9working consumers? What if all the consumers exceed timeout and are all considered dead, am I left with a running application which does nothing because there are no consumers? - Long-processing consumer finishes processing after re-balancing already took place, does broker initializes re-balance again and consumer is assigned a partition anew? As I understand it continues running #1-5 in a loop and sending a
heartbeatto broker initializes also process of adding consumer to the consumers group, from which it was removed after being givendeadstatus, correct? - Application throws some sort of exception indicating that
session.timeout.mswas exceeded and the processing is abruptly stopped? - Also what about
max.poll.interval.msproperty, what if we even exceed that period and consumer X finishes processing aftermax.poll.interval.msvalue? Consumer already exceeded thesession.timeout.msvalue, it was excluded from consumer group, status set todead, what difference does it gives us in configuring Kafka consumer?
We have a process which extracts data for processing and this extraction consists of 50+ SQL queries (majority being SELECT's, few UPDATES), they usually go fast but of course all depends on the db load and possible locks etc. and there is a possibility that the processing takes longer than the session's timeout. I do not want to infinitely increase sessions timeout until "I hit the spot". The process is idempotent, if it's repeated X times withing X minutes we do not care.