पाठ 18 / 25

Idempotent Producers और Transactions

समझें कि Kafka में exactly-once का क्या मतलब है और उसकी सीमाएँ क्या हैं।

Kafka के भीतर exactly-once

Kafka के transactions producer को कई partitions में लिखने और consumer offsets commit करने की अणु-क्रिया देते हैं: या तो सब दिखाई देता है या कुछ नहीं। इसलिए consume-process-produce ऐप topic A से पढ़ सकता है और topic B में Kafka के भीतर exactly-once semantics के साथ लिख सकता है: transactional.id तय करें, काम को begin/commit में लपेटें, और downstream consumers isolation.level=read_committed उपयोग करें ताकि वे abort हुआ डेटा कभी न देखें। गारंटी Kafka की reads और writes पर लागू है; बाहरी database भी अपडेट करते हों तो भी idempotent write या transactional outbox चाहिए, क्योंकि Kafka उन systems को roll back नहीं कर सकता जिन्हें वह नियंत्रित नहीं करता। Kafka Streams और Flink जैसे frameworks समर्थित sinks के लिए end-to-end exactly-once दे सकते हैं।

Transactional producer (उदाहरण)

orders से पढ़ता है, orders-enriched में लिखता है और offsets एक transaction में commit करता है। enrich आपके अपने function का प्रतीक है।

p = Producer({"bootstrap.servers": "localhost:9092", "transactional.id": "enricher-1"})
c = Consumer({"bootstrap.servers": "localhost:9092", "group.id": "enricher",
              "enable.auto.commit": False, "isolation.level": "read_committed"})
p.init_transactions()
c.subscribe(["orders"])
while True:
    msg = c.poll(1.0)
    if msg is None or msg.error(): continue
    p.begin_transaction()
    p.produce("orders-enriched", key=msg.key(), value=enrich(msg.value()))
    p.send_offsets_to_transaction(c.position(c.assignment()), c.consumer_group_metadata())
    p.commit_transaction()

त्वरित जाँच: `isolation.level=read_committed` consumer के लिए क्या करता है?

  • Aborted या खुले transactions के records छिपाता है
  • Records encrypt करता है
  • पुराने records हटाता है
  • Broker तेज़ करता है
Answer

Aborted या खुले transactions के records छिपाता है — Consumer को सिर्फ़ committed transactional डेटा पहुँचाया जाता है।