Consumer'ın ayrıldığını tespit etmek için iki tane yöntem var. Açıklaması şöyle.
If you don’t already know, Kafka ecosystem uses two complementary methods to spot inactive consumers in a consumer group (we don’t want a partition being kept forever by a ghost consumer):— A hearbeat thread that send periodic heartbeats to indicate that the consumer is still alive.— A proof of work progression based timeout. To sum up, you have to call poll() to prove that you’re still doing something.
1. Heartbeat Thread
Açıklaması şöyle
Every consumer is supposed to send a heartbeat periodically to the Kafka cluster. We can set that period by setting the value of a consumer property (heartbeat.interval.ms). That is used to notify the Kafka cluster that the particular consumer is alive. If a heartbeat request is not received within the configured time period (can be configured using session.timeout.ms property), the consumer will be removed from the group and coordinator will enforce a rebalance.
Yani Heartbeat yöntemini kontrol etmek için iki tane ayar var. Eğer consumer session.timeout.ms süresince heartbeat göndermezse Group Coordinator bu consumer'ı öldü kabul eder. Açıklaması şöyle
If the consumer stops sending heartbeats for long enough(session timeout session.timeout.ms), the session will time out and the group coordinator will consider it dead and trigger a rebalance. If a consumer crashes and stops processing messages, it will take the group coordinator a few seconds without heartbeats to decide it is dead and trigger the rebalance. During those seconds, no messages will be processed from the partitions owned by the dead consumer. When closing a consumer cleanly, the consumer will notify the group coordinator that it is leaving, and the group coordinator will trigger a rebalance immediately, reducing the gap in processing.
2. Poll Thread
Açıklaması şöyle
Usually the heartbeat is sent from a thread of the application. But there is a separate thread which consumes records from the Kafka topic, process (do certain operations based on record values) in the application, and again fetch another set of records from the topic using poll method. After that it goes through the same cycle over and over again.Even though the heartbeat thread is alive, consumer thread might be hanging or dead due to multiple reasons. ( Database query, network call, insufficient memory or any reason might cause this consumer thread to stop working)If that kind of thing would happen, Kafka topic coordinator should know that and remove the particular consumer (and do another rebalancing). Otherwise the partition(s) being consumed by the hanged (frozen) consumer will not be processed anymore.To know whether the consumer thread is hanging, another mechanism is is used. Consumer thread should poll records within a given time range (it can be configured with consumer property max.poll.interval.ms). If the consumer did not make a poll request within this time period, that consumer will be removed from the group and another rebalance will be occurred.
max.poll.interval.ms ise çekilen kayıtları işlemek için consumer'a tanınan süredir. Açıklaması şöyle
If a consumer doesn’t poll for a long time(max.poll.interval.ms), the heartbeat thread will stop sending heartbeats and it will send a “leave group request” to trigger a rebalance. That is good for detecting slow consumers. For example, if you poll 500 records and can’t process them in max.poll.interval.ms your consumer will leave the group and will trigger a rebalance.
Kafka yavaş Consumer'ları sevmiyor. Eğer Consumer üzün müddet çalışmıyorsa birikme olur ve poll() işlemi 500 kayıt dönebilir.
Her poll işleminde çekilecek en fazla kayıt sayısı max.poll.records ile ayarlanabilir
Yavaş consumer için hata mesajı şöyle
Consumer clientId=status-listener, groupId=status-groupId] Member my-member sending LeaveGroup request to coordinator hostname.net:9093 (id: 2147483646 rack: null) due to consumer poll timeout has expired. This means the time between subsequent calls to poll() was longer than the configured max.poll.interval.ms, which typically implies that the poll loop is spending too much time processing messages. You can address this either by increasing max.poll.interval.ms or by reducing the maximum size of batches returned in poll() with max.poll.records.
Örnek
Burada yavaş bir consumer sorusu var.
Interviewer: "Your Kafka consumers are up. No errors. The producer is sending at the same rate as always. But the lag keeps growing. Why?"Candidate: "The consumers are too slow. I would add a few more to the group."This is not answer, this is reflex.You didn't ask yourself why the lag is building in first place.To answer it correctly first you need to know a consumer is watched by the broker using 2 different clocks (settings).The first is the heartbeat. A background thread sends one every 3 seconds. Miss them for 45 seconds and the consumer is treated as dead.The second is the poll interval. It measures the gap between two calls for records. The default is 5 minutes.The heartbeat proves your consumer is alive. The poll proves your consumer is working.Now the arithmetic.Consumer polls the Kafka broker and gets 500 records by default. Assume consumer's consumption logic in this case is to call an external service that has always answered in 200 ms.So that one single consumer can process 500 records in under 2 minutes.This is working fine for a year.Today that external service is degraded and takes 2 sec to answer.Same code. Same 500 records. 16 minutes.And your consumer logic does not ask for the next batch until it has finished this one.At minute 5 the poll timer expires. Your consumer does not wait to be thrown out. It takes itself out of the group so somebody else can pick up its partitions.They go to another consumer. Nothing was committed, so this new consumer fetches the same 500 records, hits the same slow service, and resigns at minute 5 too.Meanwhile the first consumer finishes all 500 and tries to commit. It fails. It gave those partitions away 11 minutes ago.The work was done twice. The offset never moved. The lag never comes down.Now go looking for the error. There isn't one.Nothing failed, so nothing threw. A background thread logged a warning and left the group on purpose. The only honest signal is that failed commit, and it reads like a routine rebalance.Processing never reports a problem, because processing is working. It is just working on the same 500 records forever.A classic consumer group has no delivery-attempt limit and no dead letter queue. Nothing here stops on its own.You can apply these 3 fixes. Pick by what you know about the work.- Raise the poll interval to cover your real worst case. A dead consumer then takes longer to spot.- Take fewer records per poll. 500 down to 20, and even a slow day fits inside the window.- Keep polling on the main thread with the partition paused, and do the work elsewhere.The default configuration assumes your work finishes in 5 minutes. Your slowest dependency decides whether that is true.Adding consumers fixes a throughput problem. The producer never gave you one.
No comments:
Post a Comment