FastAPI बैकएंड डेवलपमेंट बूटकैंप · पाठ

Kafka इवेंट का अतुल्यकालिक उत्पादन और उपभोग

इवेंट लूप को अवरुद्ध किए बिना इवेंट प्रकाशित और उपभोग करने के लिए aiokafka को FastAPI के साथ एकीकृत कीजिए।

पाठ 1, कुल 4 में से13 चरण

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 offset

Concurrency और 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 बैकएंड डेवलपमेंट बूटकैंप पाठ में एक अंतर्निर्मित कोड संपादक शामिल है, जिससे आप सीधे अपने ब्राउज़र में वास्तविक कोड लिख और चला सकते हैं और तुरंत एआई प्रतिक्रिया पा सकते हैं—स्थानीय सेटअप की आवश्यकता नहीं है।

इस पाठ्यक्रम के सभी पाठ

  1. Kafka इवेंट का अतुल्यकालिक उत्पादन और उपभोग
  2. स्कीमा रजिस्ट्री और Avro अनुबंध विकास
  3. लेन-देनात्मक आउटबॉक्स पैटर्न
  4. इडेम्पोटेंट उपभोक्ता और ठीक-एक-बार अर्थविज्ञान
← FastAPI बैकएंड डेवलपमेंट बूटकैंप पर वापस जाएँ