लेन-देनात्मक आउटबॉक्स पैटर्न
आउटबॉक्स तालिका और रिले प्रक्रिया से स्थिति परिवर्तनों तथा इवेंट प्रकाशन की परमाण्विकता सुनिश्चित कीजिए।
लेन-देनात्मक आउटबॉक्स पैटर्न, CoddyKit पर FastAPI बैकएंड डेवलपमेंट बूटकैंप का एक निःशुल्क पाठ है। यह 4 में से 3वाँ पाठ है। इस अध्ययन पथ के 3 तक कोई भी पाठ पूरा पढ़ना निःशुल्क है — इसके बाद CoddyKit PRO हर पाठ अनलॉक करता है, साथ ही अंतर्निर्मित कोड संपादक और चौबीसों घंटे एआई शिक्षक के साथ व्यावहारिक अभ्यास भी उपलब्ध कराता है। यह FastAPI बैकएंड डेवलपमेंट बूटकैंप सीखने के मार्ग का हिस्सा है और आपकी प्रगति वेब तथा CoddyKit ऐप पर सिंक होती रहती है। FastAPI बैकएंड डेवलपमेंट बूटकैंप पाठ्यक्रम में कुल 4 पाठ शामिल हैं।
दोहरी-लेखन की समस्या
इवेंट-आधारित FastAPI सेवा में, एक अनुरोध को अक्सर दो काम करने पड़ते हैं: अपने डेटाबेस में स्थिति को स्थायी रूप से सहेजना और Kafka या Pulsar पर एक इवेंट प्रकाशित करना।
समस्या यह है कि ये दो अलग-अलग सिस्टम हैं और इनके बीच कोई साझा लेन-देन नहीं है। यदि आप DB पंक्ति को कमिट करने के बाद प्रकाशित करने से पहले क्रैश हो जाते हैं, तो डाउनस्ट्रीम सेवाओं को इसकी जानकारी कभी नहीं मिलेगी। यदि आप पहले प्रकाशित करते हैं और फिर DB कमिट विफल हो जाता है, तो आप ऐसी स्थिति के लिए इवेंट भेज देते हैं जो अस्तित्व में ही नहीं है।
db.commit()सफल,producer.send()विफल → इवेंट खो गयाproducer.send()सफल,db.commit()विफल → काल्पनिक इवेंट
इसी को दोहरी-लेखन समस्या कहते हैं, और try/except की कितनी भी मात्रा इसे पूरी तरह हल नहीं कर सकती।
async def create_order(db, producer, payload):
order = Order(**payload)
db.add(order)
await db.commit() # write #1: database
await producer.send(
"orders", order.as_event()
) # write #2: broker (may fail!)
return orderआप केवल पुनः प्रयास क्यों नहीं कर सकते
पहली सामान्य प्रतिक्रिया होती है: "मैं प्रकाशन को बस पुनः प्रयास लूप में लपेट दूँगा।" लेकिन पुनः प्रयास इस अंतर को समाप्त नहीं करते।
- कमिट और प्रकाशन के बीच प्रक्रिया को समाप्त किया जा सकता है (OOM, डिप्लॉयमेंट, k8s निष्कासन) — पुनः प्रयास करने के लिए कोई कोड नहीं चलता।
- ब्रोकर टाइमआउट के बाद पुनः प्रयास करने पर डुप्लिकेट बन सकते हैं, यदि पहला भेजना वास्तव में सफल हो चुका हो।
- एक बार आंशिक रूप से स्वीकृत हो जाने के बाद Kafka लेखन को वापस नहीं लिया जा सकता।
मूल समस्या यह है कि कमिट और प्रकाशित करने का इरादा परमाणु नहीं हैं। हमें प्रकाशन-इरादे को स्थिति परिवर्तन के साथ उसी डेटाबेस लेन-देन का हिस्सा बनाना होगा। Transactional Outbox पैटर्न ठीक यही करता है।
Outbox का विचार
Transactional Outbox पैटर्न आपके व्यावसायिक डेटा वाले उसी डेटाबेस में एक outbox तालिका जोड़ता है। ब्रोकर पर सीधे प्रकाशित करने के बजाय, आप स्थिति बदलने वाले उसी लेन-देन के भीतर इवेंट का वर्णन करने वाली एक पंक्ति INSERT करते हैं।
क्योंकि ऑर्डर पंक्ति और outbox पंक्ति एक ही लेन-देन में लिखी जाती हैं, इसलिए दोनों या तो कमिट होती हैं या दोनों वापस हो जाती हैं। ऐसी कोई स्थिति नहीं रहती जिसमें एक मौजूद हो और दूसरी न हो।
बाद में एक अलग रिले प्रक्रिया अप्रकाशित outbox पंक्तियों को पढ़ती है, उन्हें Kafka/Pulsar पर भेजती है और भेजे जाने के रूप में चिह्नित करती है। ब्रोकर लेखन अनुरोध पथ से अलग हो जाता है।
- परमाणुता वितरित लेन-देन से नहीं, बल्कि DB लेन-देन से आती है।
- डिलीवरी कम-से-कम-एक-बार हो जाती है — उपभोक्ताओं को इडेम्पोटेंट होना चाहिए।
Outbox तालिका का डिज़ाइन
Outbox तालिका में इतना मेटाडेटा होना चाहिए कि रिले सही ढंग से प्रकाशित कर सके और आप डीबग कर सकें। एक सामान्य स्कीमा इस प्रकार है:
id— UUID प्राथमिक कुंजी, जिसे इडेम्पोटेंसी के लिए इवेंट ID के रूप में भी पुनः उपयोग किया जाता हैaggregate_type/aggregate_id— इवेंट किस बारे में है (जैसेorder/ ऑर्डर ID), अक्सर पार्टिशन कुंजी के रूप में उपयोग किया जाता हैevent_type— जैसेOrderCreatedpayload— इवेंट की JSON बॉडीcreated_at— क्रम निर्धारणpublished_at— रिले द्वारा भेजे जाने तक NULL
aggregate_id को बनाए रखने से रिले Kafka पार्टिशन कुंजी निर्धारित कर सकता है, जिससे एक ऑर्डर के सभी इवेंट क्रम में बने रहते हैं।
from sqlalchemy import Column, String, DateTime, JSON, func
from sqlalchemy.dialects.postgresql import UUID
from .db import Base
import uuid
class Outbox(Base):
__tablename__ = "outbox"
id = Column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4)
aggregate_type = Column(String, nullable=False)
aggregate_id = Column(String, nullable=False)
event_type = Column(String, nullable=False)
payload = Column(JSON, nullable=False)
created_at = Column(DateTime(timezone=True), server_default=func.now())
published_at = Column(DateTime(timezone=True), nullable=True)स्थिति और इवेंट को परमाणु रूप से लिखना
यह पैटर्न का मूल आधार है। ऑर्डर और आउटबॉक्स पंक्ति को एक ही सत्र में जोड़ा और साथ में कमिट किया जाता है। ध्यान दें कि अनुरोध हैंडलर में ब्रोकर को कोई कॉल बिल्कुल नहीं है।
यदि commit() विफल होता है, तो कोई भी पंक्ति मौजूद नहीं रहती। यदि यह सफल होता है, तो दोनों मौजूद रहती हैं। इवेंट को रिले द्वारा पहुँचाया जाना सुनिश्चित है।
import uuid
async def create_order(db, payload):
order = Order(**payload)
db.add(order)
event = Outbox(
id=uuid.uuid4(),
aggregate_type="order",
aggregate_id=str(order.id),
event_type="OrderCreated",
payload={"order_id": str(order.id), **payload},
)
db.add(event)
await db.commit() # ONE transaction: state + event intent
return orderरिले: पोलिंग प्रकाशक
सबसे सरल रिले एक पोलिंग प्रकाशक है: एक पृष्ठभूमि लूप, जो बार-बार अप्रकाशित आउटबॉक्स पंक्तियों का चयन करता है, उन्हें ब्रोकर को भेजता है और प्रकाशित के रूप में चिह्नित करता है।
- वे पंक्तियाँ चुनें जहाँ
published_at IS NULLहो और उन्हेंcreated_atके अनुसार क्रमबद्ध करें। - हर पंक्ति को कुंजी के रूप में
aggregate_idका उपयोग करके Kafka/Pulsar पर प्रकाशित करें। published_at = now()सेट करें और कमिट करें।
यदि प्रक्रिया प्रकाशित करने के बाद लेकिन चिह्नित करने से पहले क्रैश हो जाती है, तो अगली पोलिंग पर पंक्ति फिर से भेजी जाती है — इसलिए कम-से-कम-एक-बार डिलीवरी होती है और idempotent उपभोक्ताओं की आवश्यकता पड़ती है। इवेंट का id डीडुप्लिकेशन कुंजी है।
async def relay_once(db, producer, batch=100):
rows = await db.fetch(
"SELECT id, aggregate_id, event_type, payload "
"FROM outbox WHERE published_at IS NULL "
"ORDER BY created_at LIMIT $1",
batch,
)
for r in rows:
await producer.send(
topic=r["event_type"],
key=r["aggregate_id"].encode(),
value=r["payload"],
headers=[("event_id", str(r["id"]).encode())],
)
await db.execute(
"UPDATE outbox SET published_at = now() WHERE id = $1",
r["id"],
)FOR UPDATE SKIP LOCKED से दोहरी प्रोसेसिंग से बचना
यदि आप थ्रूपुट बढ़ाने के लिए एक से अधिक रिले प्रतिकृतियाँ चलाते हैं, तो दो worker एक ही आउटबॉक्स पंक्तियाँ ले सकते हैं और उन्हें दो बार प्रकाशित कर सकते हैं। Postgres इसका साफ़ समाधान देता है: SELECT ... FOR UPDATE SKIP LOCKED।
FOR UPDATEचुनी गई पंक्तियों को लेन-देन के लिए लॉक करता है।SKIP LOCKEDअन्य worker को पहले से लॉक की गई पंक्तियों पर रुकने के बजाय उन्हें छोड़ने देता है।
हर worker अलग-अलग पंक्तियों का एक बैच लेता है, उन्हें प्रकाशित करता है, चिह्नित करता है और कमिट करता है — इससे लॉक खुल जाते हैं। इस तरह आउटबॉक्स किसी बाहरी कतार प्रणाली के बिना सुरक्षित समवर्ती कार्य-कतार बन जाता है।
SELECT id, aggregate_id, event_type, payload
FROM outbox
WHERE published_at IS NULL
ORDER BY created_at
LIMIT 100
FOR UPDATE SKIP LOCKED;Idempotent उपभोक्ता अनिवार्य हैं
क्योंकि आउटबॉक्स कम-से-कम-एक-बार डिलीवरी की गारंटी देता है, इसलिए हर उपभोक्ता को एक ही इवेंट को एक से अधिक बार देखने की स्थिति सहन करनी चाहिए। मानक तकनीक इवेंट आईडी द्वारा कुंजीबद्ध प्रोसेस किए गए इवेंट की तालिका है।
इवेंट लागू करने से पहले उसकी आईडी डालने का प्रयास करें। यदि वह पहले से मौजूद है, तो आप उसे देख चुके हैं — इसे छोड़ दें। डीडुप्लिकेशन प्रविष्टि और व्यावसायिक प्रभाव को एक ही लेन-देन में करें, ताकि क्रैश होने पर दोनों में असंगति न आए।
def handle_event(conn, event_id, body):
cur = conn.cursor()
try:
cur.execute(
"INSERT INTO processed_events(event_id) VALUES (%s)",
(event_id,),
)
except UniqueViolation:
conn.rollback() # duplicate -> already handled
return
apply_business_change(cur, body)
conn.commit() # dedup + effect: one transactionरिले के रूप में Change Data Capture
पोलिंग सरल है, लेकिन इससे विलंबता और DB पर भार बढ़ता है। अधिक प्रदर्शन वाला विकल्प Postgres के write-ahead log (WAL) को सीधे Change Data Capture के माध्यम से पढ़ना है — आम तौर पर Debezium का उपयोग करके।
- Debezium WAL को लगातार पढ़ता है और आउटबॉक्स तालिका में होने वाले हर INSERT के लिए Kafka संदेश भेजता है।
published_atकॉलम या पोलिंग लूप की आवश्यकता नहीं होती — लॉग ही कतार है।- इसका Outbox Event Router SMT आउटबॉक्स कॉलम को विषय, कुंजी और पेलोड से मैप करता है।
समझौता यह है कि CDC संचालन की दृष्टि से अधिक भारी है (कनेक्टर, प्रतिकृति स्लॉट), लेकिन यह लगभग रीयल-टाइम डिलीवरी देता है और आपके FastAPI ऐप से रिले का पूरा भार हटा देता है।
क्रम और विभाजन कुंजियाँ
एक ही एग्रीगेट के इवेंट आम तौर पर क्रम से पहुँचने चाहिए — OrderCreated से पहले OrderShipped। यदि आप सावधानी बरतें, तो आउटबॉक्स और ब्रोकर दोनों इसे बनाए रखते हैं:
- रिले के चयन को
created_at(या एक लगातार बढ़ने वाले अनुक्रम) के अनुसार क्रमबद्ध करें। - Kafka की विभाजन कुंजी के रूप में
aggregate_idका उपयोग करें, ताकि एक ऑर्डर के सभी इवेंट एक ही विभाजन में जाएँ, जिन्हें Kafka क्रम से पहुँचाता है।
अलग-अलग एग्रीगेट के इवेंट स्वतंत्र रूप से एक-दूसरे के बीच आ सकते हैं — आपको केवल एग्रीगेट-स्तरीय क्रम चाहिए, जो कुंजी के आधार पर विभाजन से मिलता है। सभी ऑर्डर के लिए पूर्ण वैश्विक क्रम बहुत कम आवश्यक होता है और समानांतरता समाप्त कर देता है।
# Stable hash -> partition keeps one aggregate on one partition
def partition_for(aggregate_id: str, num_partitions: int) -> int:
h = 0
for ch in aggregate_id:
h = (h * 31 + ord(ch)) & 0xFFFFFFFF
return h % num_partitions
if __name__ == "__main__":
ids = ["order-1", "order-1", "order-2", "order-3"]
for a in ids:
print(a, "->", partition_for(a, 6))उत्पादन में आउटबॉक्स का संचालन
कुछ प्रक्रियाएँ बड़े पैमाने पर आउटबॉक्स को स्वस्थ बनाए रखती हैं:
- प्रकाशित पंक्तियों को हटाएँ — समय-समय पर उन पंक्तियों को हटाएँ जिनका
published_atअवधारण अवधि से पुराना है, या उन्हें संग्रह में स्थानांतरित करें, ताकि तालिका छोटी रहे और आंशिक इंडेक्स तेज़ बना रहे। - अप्रकाशित पंक्तियों पर आंशिक इंडेक्स:
CREATE INDEX ... ON outbox (created_at) WHERE published_at IS NULL;रिले की बार-बार चलने वाली क्वेरी को कम लागत वाला बनाए रखता है। - विलंब की निगरानी करें — अप्रकाशित पंक्तियों की संख्या या आयु पर चेतावनी दें; बढ़ता हुआ बैकलॉग बताता है कि रिले बंद है या ब्रोकर तक पहुँचा नहीं जा सकता।
- Idempotent उत्पादक — ब्रोकर स्तर पर रिले द्वारा दोबारा भेजे गए डुप्लिकेट को दबाने के लिए Kafka में
enable.idempotence=trueसक्षम करें।
त्वरित जाँच: आउटबॉक्स क्यों काम करता है
यह जाँचें कि आप समझते हैं या नहीं कि Transactional Outbox पैटर्न दोहरी-लेखन समस्या का वास्तविक समाधान कैसे करता है।
पुनरावलोकन
आपने FastAPI सेवा से विश्वसनीय रूप से इवेंट प्रकाशित करना सीखा:
- दोहरी-लेखन समस्या: DB में कमिट करना और Kafka/Pulsar पर प्रकाशित करना अलग, गैर-परमाण्विक क्रियाएँ हैं, जो क्रैश होने पर असंगत हो सकती हैं।
- Transactional Outbox स्थिति परिवर्तन के साथ उसी लेन-देन में आउटबॉक्स तालिका में एक इवेंट पंक्ति लिखता है, जिससे दोनों परमाण्विक बन जाते हैं।
- एक रिले (पोलिंग प्रकाशक या Debezium CDC) बाद में आउटबॉक्स पंक्तियों को ब्रोकर को भेजता है और उन्हें प्रकाशित के रूप में चिह्नित करता है।
- सुरक्षित समवर्ती रिले के लिए
FOR UPDATE SKIP LOCKED, क्रम बनाए रखने के लिए विभाजन कुंजी के रूप मेंaggregate_id, और idempotent, कम-से-कम-एक-बार उपभोक्ताओं के लिए प्रोसेस किए गए इवेंट की तालिका का उपयोग करें। - इसे पंक्तियों की सफ़ाई, आंशिक इंडेक्स और बैकलॉग की निगरानी के साथ संचालित करें।
एआई शिक्षक के साथ FastAPI बैकएंड डेवलपमेंट बूटकैंप सीखें — निःशुल्क
अपने ब्राउज़र में वास्तविक कोड लिखें और चलाएँ, चौबीसों घंटे एआई शिक्षक से तुरंत सहायता पाएँ, और वेब या ऐप पर वहीं से शुरू करें जहाँ आपने छोड़ा था।
- पाठ्यक्रम
- 21
- पाठ
- 84
अक्सर पूछे जाने वाले प्रश्न
क्या “लेन-देनात्मक आउटबॉक्स पैटर्न” पाठ निःशुल्क है?
हाँ — FastAPI बैकएंड डेवलपमेंट बूटकैंप अध्ययन पथ के 3 तक कोई भी पाठ, जिसमें “लेन-देनात्मक आउटबॉक्स पैटर्न” भी शामिल है, यहाँ वेब पर पूरा पढ़ना निःशुल्क है। इसके बाद CoddyKit PRO हर पाठ अनलॉक करता है, साथ ही अंतर्निर्मित कोड संपादक और चौबीसों घंटे एआई शिक्षक के साथ इंटरैक्टिव अभ्यास भी उपलब्ध कराता है। FastAPI बैकएंड डेवलपमेंट बूटकैंप पाठ्यक्रम में कुल 4 पाठ शामिल हैं।
“लेन-देनात्मक आउटबॉक्स पैटर्न” में मैं क्या सीखूँगा?
आउटबॉक्स तालिका और रिले प्रक्रिया से स्थिति परिवर्तनों तथा इवेंट प्रकाशन की परमाण्विकता सुनिश्चित कीजिए। आप ब्राउज़र में सीधे चलाए जाने वाले व्यावहारिक कोड के साथ FastAPI बैकएंड डेवलपमेंट बूटकैंप का अभ्यास करते हैं, और पाठ पूरा करते समय 24/7 एआई ट्यूटर आपके प्रश्नों के उत्तर देता है।
क्या FastAPI बैकएंड डेवलपमेंट बूटकैंप शुरू करने के लिए मुझे किसी अनुभव की आवश्यकता है?
पहले के अनुभव की आवश्यकता नहीं है। CoddyKit पर FastAPI बैकएंड डेवलपमेंट बूटकैंप शुरुआती से लेकर उन्नत शिक्षार्थियों तक सभी के लिए व्यवस्थित किया गया है, इसलिए आप यहीं से या शुरुआत से सीखना शुरू कर सकते हैं और अपनी गति से आगे बढ़ सकते हैं। यह 4 में से 3वाँ पाठ है।
“लेन-देनात्मक आउटबॉक्स पैटर्न” पाठ पूरा करने में कितना समय लगता है?
CoddyKit का अधिकांश पाठ लगभग 5–10 मिनट में पूरा हो जाता है। हर पाठ छोटा और संवादात्मक है, इसलिए आप लगातार प्रगति करते हैं और वेब या ऐप पर वहीं से सीखना जारी रख सकते हैं जहाँ आपने छोड़ा था।
क्या मैं इस FastAPI बैकएंड डेवलपमेंट बूटकैंप पाठ में कोड लिख और चला सकता हूँ?
हाँ। हर FastAPI बैकएंड डेवलपमेंट बूटकैंप पाठ में एक अंतर्निर्मित कोड संपादक शामिल है, जिससे आप सीधे अपने ब्राउज़र में वास्तविक कोड लिख और चला सकते हैं और तुरंत एआई प्रतिक्रिया पा सकते हैं—स्थानीय सेटअप की आवश्यकता नहीं है।
इस पाठ्यक्रम के सभी पाठ
- Kafka इवेंट का अतुल्यकालिक उत्पादन और उपभोग
- स्कीमा रजिस्ट्री और Avro अनुबंध विकास
- लेन-देनात्मक आउटबॉक्स पैटर्न
- इडेम्पोटेंट उपभोक्ता और ठीक-एक-बार अर्थविज्ञान