पाठ 9 / 25

Consumer Groups और Committed Offsets

समझाएँ कि group partitions कैसे साझा करता है और committed offsets व lag कैसे ट्रैक होते हैं।

प्रति group, एक partition पर एक consumer

जो consumers एक ही group.id साझा करते हैं वे consumer group बनाते हैं। Kafka हर partition को group में ठीक एक consumer को देता है, इसलिए group topic को समानांतर प्रोसेस करता है, जबकि अलग group पूरा topic स्वतंत्र रूप से पढ़ता है। काम करते हुए consumer अगले पढ़े जाने वाले record का offset commit करता है, जो आंतरिक topic (__consumer_offsets) में रखा जाता है। Lag प्रति partition log-end offset घटा committed offset है: बढ़ता lag मतलब consumers पिछड़ रहे हैं। Group partitions से ज़्यादा consumers उपयोग नहीं कर सकता; अतिरिक्त consumers निष्क्रिय रहते हैं।

काम बाँटें, स्थिति ट्रैक करें

Group के consumers topic के partitions बाँटते हैं; Kafka हर group के committed offsets याद रखता है।

चार विचार: group, assignment, offset, rebalance।
चित्र 3.1 — Group, assignment, offset और rebalance।

असली broker से group lag, चलाकर

मैंने group billing से 7 में से 4 records पढ़े, फिर kafka-consumer-groups.sh --describe चलाया। Partitions 1 और 2 पूरी तरह पढ़े गए (lag 0); partitions 4 और 5 में अब भी 1 और 2 अपठित records हैं (कुल lag 3)। Group में कोई सक्रिय सदस्य नहीं क्योंकि console consumer बाहर निकल गया।

kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic orders \
  --group billing --from-beginning --max-messages 4
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group billing

Output:

Consumer group 'billing' has no active members.

GROUP   TOPIC   PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
billing orders  0          0               0               0
billing orders  1          1               1               0
billing orders  2          3               3               0
billing orders  3          0               0               0
billing orders  4          0               1               1
billing orders  5          0               2               2

Lag का गणित, चलाकर

कुल lag partitions का योग है; यहाँ 0 + 80 + 420 = 500 records पीछे।

log_end   = {0: 1500, 1: 1480, 2: 1520}
committed = {0: 1500, 1: 1400, 2: 1100}
lag = {p: log_end[p] - committed[p] for p in log_end}
print(lag, sum(lag.values()))

Output:

{0: 0, 1: 80, 2: 420} 500

त्वरित जाँच: Consumer lag क्या है?

  • Broker तक network latency
  • Topic की उम्र
  • Log-end offset घटा committed offset
  • Partitions की संख्या
Answer

Log-end offset घटा committed offset — यह गिनता है कि कितने records लिखे जा चुके हैं पर group ने अब तक प्रोसेस नहीं किए।