पाठ 10 / 25

Assignment और Rebalancing

देखें कि partitions consumers में कैसे बँटते हैं और सदस्यों के जुड़ने या जाने पर क्या होता है।

सदस्यता बदलने पर partitions चलते हैं

जब कोई consumer जुड़ता है, जाता है, crash होता है (session.timeout.ms के भीतर heartbeats चूकता है) या polls के बीच बहुत समय लेता है (max.poll.interval.ms), तो group rebalance करता है: partitions फिर से बाँटे जाते हैं। Assignor तय करता है कैसे: range और roundrobin पारंपरिक हैं; cooperative-sticky ज़्यादातर assignments जगह पर रखता है और सिर्फ़ ज़रूरी चलाता है, पूरे group को रोकने से बचाता है। बार-बार rebalance throughput घटाते हैं, इसलिए प्रति poll प्रोसेसिंग छोटी रखें, timeouts समझदारी से तय करें और cooperative assignment चुनें। 6-partition topic में 2 consumers को 3-3 partitions, 6 consumers को एक-एक मिलते हैं, और 7वें को कोई नहीं।

Range assignment, चलाकर

मैंने range रणनीति का मॉडल चलाया। दो consumers को 3-3 partitions मिलते हैं; सात consumers में सातवाँ निष्क्रिय रहता है; एक consumer सारे छह लेता है।

def range_assign(parts, consumers):
    out = {c: [] for c in consumers}
    per, extra = divmod(len(parts), len(consumers)); i = 0
    for idx, c in enumerate(consumers):
        take = per + (1 if idx < extra else 0)
        out[c] = parts[i:i + take]; i += take
    return out

print(range_assign(list(range(6)), ["c1", "c2"]))
print(range_assign(list(range(6)), ["c1", "c2", "c3", "c4", "c5", "c6", "c7"]))
print(range_assign(list(range(6)), ["c1"]))

Output:

{'c1': [0, 1, 2], 'c2': [3, 4, 5]}
{'c1': [0], 'c2': [1], 'c3': [2], 'c4': [3], 'c5': [4], 'c6': [5], 'c7': []}
{'c1': [0, 1, 2, 3, 4, 5]}

महत्वपूर्ण consumer settings (उदाहरण)

स्थिर group के लिए आम properties।

group.id=billing
auto.offset.reset=earliest
enable.auto.commit=false
max.poll.records=200
max.poll.interval.ms=300000
session.timeout.ms=30000
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor

त्वरित जाँच: Topic में 6 partitions हैं और group में 8 consumers। क्या होगा?

  • 6 consumers को एक-एक partition मिलता है और 2 निष्क्रिय रहते हैं
  • सभी 8 बराबर partitions साझा करते हैं
  • Group शुरू नहीं होता
  • Kafka 2 नए partitions अपने आप बनाता है
Answer

6 consumers को एक-एक partition मिलता है और 2 निष्क्रिय रहते हैं — एक partition प्रति group अधिकतम एक consumer को जाता है, इसलिए अतिरिक्त consumers निष्क्रिय हैं।