Kafka इवेंट का अतुल्यकालिक उत्पादन और उपभोग
इवेंट लूप को अवरुद्ध किए बिना इवेंट प्रकाशित और उपभोग करने के लिए aiokafka को FastAPI के साथ एकीकृत कीजिए।
Kafka इवेंट का अतुल्यकालिक उत्पादन और उपभोग, CoddyKit पर FastAPI बैकएंड डेवलपमेंट बूटकैंप का एक निःशुल्क पाठ है। यह 4 में से 1वाँ पाठ है। इस अध्ययन पथ के 3 तक कोई भी पाठ पूरा पढ़ना निःशुल्क है — इसके बाद CoddyKit PRO हर पाठ अनलॉक करता है, साथ ही अंतर्निर्मित कोड संपादक और चौबीसों घंटे एआई शिक्षक के साथ व्यावहारिक अभ्यास भी उपलब्ध कराता है। यह FastAPI बैकएंड डेवलपमेंट बूटकैंप सीखने के मार्ग का हिस्सा है और आपकी प्रगति वेब तथा CoddyKit ऐप पर सिंक होती रहती है। FastAPI बैकएंड डेवलपमेंट बूटकैंप पाठ्यक्रम में कुल 4 पाठ शामिल हैं।
FastAPI में async Kafka क्यों
FastAPI asyncio इवेंट लूप पर चलता है। यदि आप Kafka को किसी ब्लॉक करने वाले क्लाइंट (जैसे मानक kafka-python लाइब्रेरी) से प्रकाशित या पोल करते हैं, तो हर नेटवर्क कॉल पूरे लूप को रोक देती है और सभी समवर्ती अनुरोध ठहर जाते हैं।
- aiokafka एक मूल asyncio Kafka क्लाइंट है, जो लूप को कभी ब्लॉक नहीं करता।
- इसके
sendऔरgetoneऑपरेशन coroutines हैं, जिन पर आपawaitलगाते हैं। - इससे एक worker Kafka I/O के लंबित रहने के दौरान हज़ारों चल रहे अनुरोध संभाल सकता है।
इस पाठ में आप AIOKafkaProducer और AIOKafkaConsumer को FastAPI ऐप में सही तरीके से जोड़ेंगे।
lifespan के साथ producer का जीवनचक्र
एक producer TCP कनेक्शन और पृष्ठभूमि में चलने वाला sender task बनाए रखता है। आपको ऐप शुरू होते समय इसे एक बार start() करना चाहिए और बंद होते समय stop() करना चाहिए — हर अनुरोध पर कभी नहीं। FastAPI में आधुनिक तरीका lifespan context manager है।
await producer.start()कनेक्शन और sender loop खोलता है।await producer.stop()लंबित batches को भेजकर साफ़ तरीके से बंद करता है।- producer को
app.stateपर रखें, ताकि routes उस तक पहुँच सकें।
from contextlib import asynccontextmanager
from fastapi import FastAPI
from aiokafka import AIOKafkaProducer
@asynccontextmanager
async def lifespan(app: FastAPI):
producer = AIOKafkaProducer(
bootstrap_servers="localhost:9092",
enable_idempotence=True,
)
await producer.start()
app.state.producer = producer
try:
yield
finally:
await producer.stop()
app = FastAPI(lifespan=lifespan)किसी route से event प्रकाशित करना
किसी route के भीतर साझा producer प्राप्त करें और await producer.send_and_wait(...) चलाएँ। send_and_wait कॉल broker द्वारा record की पुष्टि किए जाने पर लौटती है, जिससे आपको back-pressure और डिलीवरी की पुष्टि मिलती है।
- Kafka keys और values bytes होते हैं — JSON को स्वयं encode करें या serializer दें।
- लौटाया गया
RecordMetadataआपको partition और offset बताता है। - partitions के बीच हर entity का क्रम सुनिश्चित करने के लिए message की key का उपयोग करें।
import json
from fastapi import FastAPI, Request
app = FastAPI()
@app.post("/orders")
async def create_order(payload: dict, request: Request):
producer = request.app.state.producer
value = json.dumps(payload).encode("utf-8")
key = str(payload["order_id"]).encode("utf-8")
meta = await producer.send_and_wait(
"orders", value=value, key=key
)
return {"partition": meta.partition, "offset": meta.offset}Serializer बनाम मैन्युअल encoding
हर send पर json.dumps(...).encode() कॉल करने के बजाय, आप aiokafka को value_serializer और key_serializer दे सकते हैं। producer उन्हें अपने-आप लागू करता है, इसलिए routes सीधे Python objects भेजते हैं।
value_serializerआपका object प्राप्त करता है और उसेbytesके रूप में लौटाना आवश्यक है।- इससे encoding एक ही जगह केंद्रित हो जाती है और अलग-अलग routes में दोहराव से बचता है।
- उत्पादन परिवेश में टीमें अक्सर schema registry के माध्यम से JSON की जगह Avro या Protobuf इस्तेमाल करती हैं।
import json
from aiokafka import AIOKafkaProducer
producer = AIOKafkaProducer(
bootstrap_servers="localhost:9092",
value_serializer=lambda v: json.dumps(v).encode("utf-8"),
key_serializer=lambda k: str(k).encode("utf-8"),
)
# Now routes can send native objects:
# await producer.send_and_wait("orders", value={"id": 7}, key=7)send बनाम send_and_wait
दो producer methods, दो प्रकार के समझौते:
send()तुरंत एक future लौटाता है और aiokafka को पृष्ठभूमि में records का batch बनाने देता है — throughput अधिक होता है, लेकिन अभी यह पता नहीं चलता कि डिलीवरी सफल हुई या नहीं।send_and_wait()broker की पुष्टि का इंतज़ार करता है — हर कॉल थोड़ी धीमी होती है, लेकिन परिणाम मिलता है और त्रुटियाँ तुरंत सामने आ जाती हैं।
जिन request handlers में client को पुष्टि चाहिए, वहाँ send_and_wait को प्राथमिकता दें। बिना पुष्टि के बड़े पैमाने पर भेजने के लिए send का उपयोग करें और चाहें तो futures पर बाद में await करें।
# Fire many records fast, then wait once for all of them
futures = []
for item in batch:
fut = await producer.send("events", value=item)
futures.append(fut)
# Awaiting the futures surfaces any delivery errors
for fut in futures:
record_meta = await futजिस blocking जाल से बचना है
सबसे सामान्य गलती async ऐप में synchronous client को मिलाना है। नीचे दिया गया स्वतंत्र demo बताता है कि लूप को block करना क्यों नुकसानदायक है: coroutine के भीतर blocking time.sleep सब कुछ रोक देता है, जबकि asyncio.sleep नियंत्रण छोड़ देता है।
इसे चलाकर देखें कि await वाला संस्करण दोनों tasks को एक-दूसरे के साथ चलने देता है — बिल्कुल वही व्यवहार जो aiokafka वास्तविक नेटवर्क I/O के लिए देता है।
import asyncio
import time
async def good_io(name):
await asyncio.sleep(0.2) # yields the loop
print(f"{name} done at {time.strftime('%X')}")
async def main():
start = time.perf_counter()
await asyncio.gather(good_io("A"), good_io("B"))
print(f"both finished in {time.perf_counter() - start:.2f}s")
asyncio.run(main())पृष्ठभूमि task के रूप में consumer
FastAPI ऐप HTTP सेवा देता है, लेकिन Kafka consumer को लगातार poll करना पड़ता है। साफ़ पैटर्न यह है: consumer को lifespan में शुरू करें और उसके poll loop को asyncio background task के रूप में चलाएँ, फिर बंद करते समय उसे रद्द करें।
asyncio.create_task(...)startup को block किए बिना loop शुरू करता है।- बंद करते समय task को रद्द करें, फिर
await consumer.stop()चलाएँ। - loop के body को हमेशा इस तरह लपेटें कि एक खराब message consumer को बंद न कर दे।
import asyncio
from contextlib import asynccontextmanager
from fastapi import FastAPI
from aiokafka import AIOKafkaConsumer
async def consume(consumer: AIOKafkaConsumer):
async for msg in consumer:
try:
handle(msg.value)
except Exception as exc:
print("handler failed:", exc)
@asynccontextmanager
async def lifespan(app: FastAPI):
consumer = AIOKafkaConsumer(
"orders",
bootstrap_servers="localhost:9092",
group_id="order-workers",
)
await consumer.start()
task = asyncio.create_task(consume(consumer))
try:
yield
finally:
task.cancel()
await consumer.stop()
app = FastAPI(lifespan=lifespan)Messages पर चलना और decoding
एक AIOKafkaConsumer async iterator है: async for msg in consumer हर नए record के लिए await करता है। हर msg में topic, partition, offset, key और value bytes के रूप में उपलब्ध होते हैं।
msg.valueको उसी तरह decode करें, जिस तरह producer की ओर से उसे encode किया था।- आप consumer constructor को
value_deserializerभी दे सकते हैं। msg.timestampमें विलंबता मेट्रिक्स के लिए broker या producer का timestamp होता है।
import json
from aiokafka import AIOKafkaConsumer
consumer = AIOKafkaConsumer(
"orders",
bootstrap_servers="localhost:9092",
group_id="order-workers",
value_deserializer=lambda b: json.loads(b.decode("utf-8")),
)
async def run():
async for msg in consumer:
event = msg.value # already a dict
print(event["order_id"], "at offset", msg.offset)Offset commit: कम-से-कम-एक बार डिलीवरी
Auto-commit (डिफ़ॉल्ट) समय-समय पर offsets commit करता है, जिससे worker के commit करने के बाद लेकिन processing से पहले crash होने पर messages खो सकते हैं। विश्वसनीय processing के लिए enable_auto_commit=False सेट करें और message को संभालना पूरा करने के बाद commit करें।
- Manual commit कम-से-कम-एक बार semantics देता है — crash होने पर आखिरी uncommitted message फिर से चलता है।
- Messages दोहराए जा सकते हैं, इसलिए आपके handlers idempotent होने चाहिए।
- Throughput और दोबारा चलाने की लागत के बीच संतुलन के लिए छोटे batches में commit करें।
consumer = AIOKafkaConsumer(
"orders",
bootstrap_servers="localhost:9092",
group_id="order-workers",
enable_auto_commit=False,
auto_offset_reset="earliest",
)
async def run():
async for msg in consumer:
await process(msg.value) # do the work first
await consumer.commit() # then advance the offsetConcurrency और partition का क्रम
एक ही partition के भीतर Kafka क्रम बनाए रखता है और async-for loop उन records को क्रमिक रूप से process करता है। विस्तार के लिए आपके पास दो विकल्प हैं:
- अधिक partitions + अधिक consumers एक ही
group_idमें — Kafka instances के बीच partitions अपने-आप बाँट देता है। - हर worker में सीमित concurrency के लिए semaphore का उपयोग करें, जब हर message का काम I/O-प्रधान हो और सख्त क्रम आवश्यक न हो।
यदि entity के अनुसार क्रम महत्वपूर्ण है, तो स्थिर message key के माध्यम से उस entity को एक ही partition में रखें और उसके partition को क्रमिक रूप से process करें।
import asyncio
sem = asyncio.Semaphore(10)
async def handle_bounded(value):
async with sem:
await do_async_work(value)
async def run(consumer):
async for msg in consumer:
# Schedule work without blocking the poll loop
asyncio.create_task(handle_bounded(msg.value))सुव्यवस्थित shutdown और error handling
एक मजबूत deployment सफ़ाई से बंद होता है, ताकि चल रहा data न खोए:
- बंद करते समय cancel consumer task करें और
await consumer.stop()चलाएँ — इससे offsets commit होते हैं (यदि auto-commit चालू है) और group साफ़ तरीके से छोड़ा जाता है, जिससे rebalancing तेज़ होती है। await producer.stop()बंद करने से पहले buffered records को भेज देता है।- हर message की handling को try/except में लपेटें और loop को crash कराने के बजाय खराब messages को dead-letter topic पर भेजें।
Request handlers के भीतर कभी start()/stop() न चलाएँ — इससे connections पर अनावश्यक दबाव पड़ता है और consumer group की सदस्यता टूटती है।
async def consume(consumer, dlq_producer):
async for msg in consumer:
try:
await process(msg.value)
await consumer.commit()
except PermanentError:
await dlq_producer.send_and_wait(
"orders.dlq", value=msg.value, key=msg.key
)
await consumer.commit() # skip the poison messageत्वरित जाँच
आपको ऐसी विश्वसनीय processing चाहिए जिसमें worker के crash होने पर order event चुपचाप कभी न खोए। कौन-सा consumer configuration इसके लिए सबसे अच्छा है?
पुनरावलोकन
आपने इवेंट लूप को अवरुद्ध किए बिना Kafka को FastAPI में एकीकृत किया है:
- aiokafka asyncio के लिए मूल रूप से बनाए गए awaitable उत्पादक और उपभोक्ता क्लाइंट प्रदान करता है।
- उत्पादक और उपभोक्ता को लाइफ़स्पैन संदर्भ में प्रबंधित करें: प्रारंभ पर
start(), बंद करते समयstop()— हर अनुरोध पर कभी नहीं। send_and_waitडिलीवरी की पुष्टि करता है;sendअधिकतम थ्रूपुट प्रदान करता है।- उपभोक्ता के पोल लूप को asyncio की बैकग्राउंड टास्क के रूप में चलाएँ और
async forसे पुनरावृत्ति करें। - ऑटो-कमिट अक्षम करें और कम-से-कम-एक-बार डिलीवरी के लिए प्रोसेसिंग के बाद कमिट करें, तथा हैंडलरों को इडेम्पोटेंट रखें।
- अधिक पार्टिशन और समूह में अधिक उपभोक्ताओं के साथ स्केल करें, संदेश कुंजी के माध्यम से प्रत्येक इकाई का क्रम बनाए रखें, और खराब संदेशों के लिए डेड-लेटर टॉपिक के साथ व्यवस्थित रूप से बंद करें।
एआई शिक्षक के साथ FastAPI बैकएंड डेवलपमेंट बूटकैंप सीखें — निःशुल्क
अपने ब्राउज़र में वास्तविक कोड लिखें और चलाएँ, चौबीसों घंटे एआई शिक्षक से तुरंत सहायता पाएँ, और वेब या ऐप पर वहीं से शुरू करें जहाँ आपने छोड़ा था।
- पाठ्यक्रम
- 21
- पाठ
- 84
अक्सर पूछे जाने वाले प्रश्न
क्या “Kafka इवेंट का अतुल्यकालिक उत्पादन और उपभोग” पाठ निःशुल्क है?
हाँ — FastAPI बैकएंड डेवलपमेंट बूटकैंप अध्ययन पथ के 3 तक कोई भी पाठ, जिसमें “Kafka इवेंट का अतुल्यकालिक उत्पादन और उपभोग” भी शामिल है, यहाँ वेब पर पूरा पढ़ना निःशुल्क है। इसके बाद CoddyKit PRO हर पाठ अनलॉक करता है, साथ ही अंतर्निर्मित कोड संपादक और चौबीसों घंटे एआई शिक्षक के साथ इंटरैक्टिव अभ्यास भी उपलब्ध कराता है। FastAPI बैकएंड डेवलपमेंट बूटकैंप पाठ्यक्रम में कुल 4 पाठ शामिल हैं।
“Kafka इवेंट का अतुल्यकालिक उत्पादन और उपभोग” में मैं क्या सीखूँगा?
इवेंट लूप को अवरुद्ध किए बिना इवेंट प्रकाशित और उपभोग करने के लिए aiokafka को FastAPI के साथ एकीकृत कीजिए। आप ब्राउज़र में सीधे चलाए जाने वाले व्यावहारिक कोड के साथ FastAPI बैकएंड डेवलपमेंट बूटकैंप का अभ्यास करते हैं, और पाठ पूरा करते समय 24/7 एआई ट्यूटर आपके प्रश्नों के उत्तर देता है।
क्या FastAPI बैकएंड डेवलपमेंट बूटकैंप शुरू करने के लिए मुझे किसी अनुभव की आवश्यकता है?
पहले के अनुभव की आवश्यकता नहीं है। CoddyKit पर FastAPI बैकएंड डेवलपमेंट बूटकैंप शुरुआती से लेकर उन्नत शिक्षार्थियों तक सभी के लिए व्यवस्थित किया गया है, इसलिए आप यहीं से या शुरुआत से सीखना शुरू कर सकते हैं और अपनी गति से आगे बढ़ सकते हैं। यह 4 में से 1वाँ पाठ है।
“Kafka इवेंट का अतुल्यकालिक उत्पादन और उपभोग” पाठ पूरा करने में कितना समय लगता है?
CoddyKit का अधिकांश पाठ लगभग 5–10 मिनट में पूरा हो जाता है। हर पाठ छोटा और संवादात्मक है, इसलिए आप लगातार प्रगति करते हैं और वेब या ऐप पर वहीं से सीखना जारी रख सकते हैं जहाँ आपने छोड़ा था।
क्या मैं इस FastAPI बैकएंड डेवलपमेंट बूटकैंप पाठ में कोड लिख और चला सकता हूँ?
हाँ। हर FastAPI बैकएंड डेवलपमेंट बूटकैंप पाठ में एक अंतर्निर्मित कोड संपादक शामिल है, जिससे आप सीधे अपने ब्राउज़र में वास्तविक कोड लिख और चला सकते हैं और तुरंत एआई प्रतिक्रिया पा सकते हैं—स्थानीय सेटअप की आवश्यकता नहीं है।
इस पाठ्यक्रम के सभी पाठ
- Kafka इवेंट का अतुल्यकालिक उत्पादन और उपभोग
- स्कीमा रजिस्ट्री और Avro अनुबंध विकास
- लेन-देनात्मक आउटबॉक्स पैटर्न
- इडेम्पोटेंट उपभोक्ता और ठीक-एक-बार अर्थविज्ञान