एक VPS पर NATS, RabbitMQ और Kafka: क्या चुनें?
एक सर्वर पर message queue चुनते समय delivery guarantees, मेमोरी लागत और restart व्यवहार को समझें। जानें कि कब Postgres का उपयोग करना Kafka से बेहतर विकल्प साबित होता है।
एक सर्वर के लिए संक्षिप्त उत्तर
एक VPS पर message queue का उपयोग करना delivery guarantees के बारे में एक निर्णय है, न कि गति के बारे में। एक सिंगल बॉक्स पर broker शायद ही कभी bottleneck बनता है, क्योंकि आपका application code, database और एक ही disk पहले ही वहां मौजूद होते हैं। उस टूल को चुनें जिसके failure behaviour के साथ आप काम करने को तैयार हैं, और फिर अपने पास मौजूद बॉक्स को मापें।
चार विकल्प, उसी क्रम में जिस पर अधिकांश पाठकों को विचार करना चाहिए।
- उस database का उपयोग करें जिसे आप पहले से चला रहे हैं। Postgres और
SELECT ... FOR UPDATE SKIP LOCKEDएक कार्यशील job queue है, और यह monitor करने के लिए कोई नई process नहीं जोड़ता है। - RabbitMQ का उपयोग तब करें जब प्रत्येक message काम की एक ऐसी इकाई हो जिसे acknowledge किया जाना चाहिए, जिसे सीमित बार retry किया जाना चाहिए, और फिर कहीं ऐसी जगह रखा जाना चाहिए जहाँ कोई इंसान उसे देख सके।
- NATS का उपयोग तब करें जब messages ऐसे events हों जिन पर आपके सिस्टम के कई हिस्से प्रतिक्रिया देते हैं। उन events के लिए JetStream चालू करें जिन्हें restart के बाद भी सुरक्षित रहना चाहिए।
- Kafka का उपयोग तब करें जब कोई downstream टूल केवल Kafka protocol का ही समर्थन करता हो। एक सर्वर पर यही एकमात्र कारण बचा है।
यह गाइड बाकी का तर्क प्रदान करती है: एक छोटे VPS पर प्रत्येक विकल्प की memory और disk लागत क्या है, बॉक्स के reboot होने पर यह क्या करता है, और वह सटीक command जो आपके उपयोगकर्ताओं को महसूस होने से पहले backlog दिखाती है।
डिलीवरी गारंटी का वास्तविक अर्थ क्या है
At most once का अर्थ है कि ब्रोकर मैसेज को आगे बढ़ा देता है और उसे भूल जाता है। यदि कोई कंज्यूमर कनेक्टेड नहीं है, या काम पूरा होने से पहले ही कंज्यूमर बंद हो जाता है, तो मैसेज खो जाता है और इसकी कोई रिपोर्ट नहीं मिलती।
At least once का अर्थ है कि काम सफलतापूर्वक पूरा होने के बाद कंज्यूमर एक एकनॉलेजमेंट (ack) भेजता है। जब तक वह ack प्राप्त नहीं होता, ब्रोकर मैसेज को अपने पास रखता है और उसे दोबारा डिलीवर करेगा। रीडिलीवरी के कारण ही आपके हैंडलर्स का 'idempotent' होना आवश्यक है: एक ही मैसेज को दो बार प्रोसेस करने पर कार्ड से दो बार पैसे नहीं कटने चाहिए। एंड-टू-एंड 'Exactly once' की सुविधा कोई ब्रोकर नहीं देता। यह आपके अपने डेटाबेस में मौजूद एक यूनिक की (unique key) से सुनिश्चित होता है।
Replay एक अलग विशेषता है। एक क्यू (queue) एकनॉलेज होने के बाद मैसेज को हटा देती है। एक लॉग (log) इसे एक निश्चित रिटेंशन विंडो तक रखता है, ताकि एक नया कंज्यूमर शुरुआत से पूरा इतिहास पढ़ सके। Kafka और NATS JetStream लॉग हैं। RabbitMQ एक क्यू है। यह अंतर थ्रूपुट (throughput) की तुलना में अधिक आर्किटेक्चर को प्रभावित करता है।
Dead lettering उस मैसेज के साथ होता है जो बार-बार फेल हो रहा हो। इसके बिना, एक 'poison message' हमेशा के लिए लूप में फंस जाता है, और यह लूप किसी खराब वर्कर के बजाय एक व्यस्त वर्कर जैसा दिखाई देता है।
Postgres के साथ शुरुआत करें और broker को खुद को साबित करने दें
अधिकांश single-application workloads में प्रतिदिन कुछ हजार background jobs होती हैं। यह एक table में आसानी से समा जाती हैं।
CREATE TABLE job (
id bigserial PRIMARY KEY,
payload jsonb NOT NULL,
run_after timestamptz NOT NULL DEFAULT now(),
attempts int NOT NULL DEFAULT 0
);
CREATE INDEX job_ready_idx ON job (run_after, id);एक worker एक transaction के भीतर एक job को claim करता है।
BEGIN;
SELECT id, payload
FROM job
WHERE run_after <= now()
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 1;
-- run the work, then remove the row
DELETE FROM job WHERE id = $1;
COMMIT;FOR UPDATE SKIP LOCKED ही पूरी तरकीब है। यह उस row को lock कर देता है जिसे यह return करता है और उन rows को छोड़ देता है जिन्हें किसी अन्य transaction ने पहले ही lock कर रखा है, ताकि दो workers कभी भी एक ही job को claim न करें। यदि कोई worker crash हो जाता है, तो Postgres उसके transaction को abort कर देता है, lock release हो जाता है, और row अगले worker के लिए उपलब्ध हो जाती है। आपको at-least-once delivery, attempts को बढ़ाकर retries, और एक dead letter table मिलती है, वह भी उस durability के साथ जिसके लिए आप पहले से भुगतान कर रहे हैं। backlog के लिए केवल एक query की आवश्यकता होती है: SELECT count(*) FROM job WHERE run_after <= now();
जहाँ यह काम करना बंद कर देता है। प्रत्येक claim और delete एक write operation है, इसलिए उच्च job rate के कारण dead row versions पीछे छूट जाते हैं, और एक queue table bloat का क्लासिक उदाहरण है जो autovacuum की गति से अधिक तेजी से बढ़ता है। लंबी jobs स्थिति को और खराब कर देती हैं, क्योंकि काम की अवधि तक खुला रहने वाला transaction पूरे database के लिए vacuum horizon को रोक कर रखता है। Polling से latency बढ़ती है, और LISTEN के साथ NOTIFY polling को तो हटा देते हैं लेकिन writes को नहीं। जब job table आपकी सबसे व्यस्त table हो, या किसी दूसरी service को उन्हीं events की आवश्यकता हो, तो काम को बाहर ले जाएँ। यह विकल्प इस बात पर निर्भर करता है कि database स्वयं कैसे deploy किया गया है, इसलिए इसके बगल में broker जोड़ने से पहले यह तय करें कि database Docker में चलता है या host पर।
Redis दूसरी ऐसी चीज़ है जिसे आप शायद पहले से चला रहे हों। Redis Streams आपको XADD और XREADGROUP के साथ consumer groups, प्रति group एक pending list, और मृत हो चुके consumer से काम वापस लेने के लिए XAUTOCLAIM प्रदान करते हैं। यह छोटा और तेज़ है। एक box पर इसकी वास्तविक सीमा यह है: सामान्य appendfsync everysec setting के साथ, बिजली गुल होने पर लगभग एक सेकंड के writes खो सकते हैं। यह cache invalidation के लिए ठीक है लेकिन payments के लिए गलत है। यदि आपका application VPS पर production में SQLite के इर्द-गिर्द बना एक single process है, तो वही claim-and-delete pattern काम करता है, हालाँकि SQLite में SKIP LOCKED का कोई समकक्ष नहीं है और प्रत्येक worker एक ही write lock पर serialize होता है।
NATS core: बिना मेमोरी वाला subject routing
docker run -d --name nats \
-p 4222:4222 -p 127.0.0.1:8222:8222 \
nats:2.14 -m 8222अगस्त 2026 तक, वर्तमान सर्वर लाइन 2.14 है। -m 8222 HTTP मॉनिटरिंग पोर्ट को चालू करता है, जो डिफ़ॉल्ट रूप से बंद रहता है और इसमें कोई प्रमाणीकरण नहीं होता है, इसलिए इसे ऊपर बताए अनुसार localhost पर ही bind करें।
Core NATS 'at most once' (अधिकतम एक बार) डिलीवरी मॉडल पर काम करता है और यह कुछ भी स्टोर नहीं करता है। एक प्रकाशक (publisher) orders.created जैसे subject पर संदेश भेजता है, और जिस भी ग्राहक (subscriber) का फ़िल्टर मेल खाता है, उसे एक कॉपी मिल जाती है। यदि कोई भी सब्सक्राइब नहीं है, तो संदेश को हटा दिया जाता है और प्रकाशक को कोई त्रुटि नहीं दिखाई देती है, क्योंकि प्रकाशक का कार्य सर्वर द्वारा बाइट्स स्वीकार करते ही समाप्त हो जाता है। एक queue group (एक ही समूह नाम साझा करने वाले कई ग्राहक) सर्वर को प्रति संदेश एक सदस्य चुनने के लिए कहता है, जो बिना किसी कतार (queue) को स्टोर किए कार्य का वितरण करता है।
इसका फुटप्रिंट केवल सब्सक्रिप्शन स्टेट और प्रत्येक कनेक्शन के लिए एक राइट बफ़र होता है, इसलिए यह संदेश की मात्रा के बजाय कनेक्शन की संख्या को ट्रैक करता है, और डिस्क पर कुछ भी जमा नहीं होता है। रीस्टार्ट का व्यवहार इसी से तय होता है: इन-फ़्लाइट संदेश समाप्त हो जाते हैं, क्लाइंट अपने आप फिर से कनेक्ट हो जाते हैं, और प्रतीक्षा करने के लिए कोई रिकवरी चरण नहीं होता है।
यहाँ निगरानी के लिए कोई बैकलॉग नहीं होता है, इसलिए संदेश हानि (loss) पर नज़र रखें। जब कोई ग्राहक सर्वर द्वारा लिखे जाने की तुलना में अपने सॉकेट को धीमी गति से पढ़ता है, तो उस क्लाइंट के लिए सर्वर का बफ़र भर जाता है। यदि क्लाइंट राइट डेडलाइन तक संदेशों को प्रोसेस नहीं कर पाता है, तो सर्वर पूरा कनेक्शन बंद कर देता है और एक काउंटर बढ़ा देता है।
curl -s http://localhost:8222/varz | jq '.slow_consumers, .connections, .in_msgs, .out_msgs'slow_consumers का मान यदि लगातार बढ़ रहा है, तो इसका मतलब है कि संदेश ड्रॉप हो रहे हैं, इसलिए इसे एक बार पढ़ने के बजाय इस पर अलर्ट सेट करें। Core NATS उन संदेशों के लिए उपयुक्त है जिनका मूल्य जल्दी समाप्त हो जाता है: जैसे कि कोई मेट्रिक, उपस्थिति अपडेट (presence update), या कैश इनवैलिडेशन जिसे अगला इवेंट वैसे भी बदल देगा।
NATS JetStream: एक ही प्रोसेस में durable streams और replay
JetStream कोई अलग उत्पाद नहीं है। यह उसी बाइनरी में मौजूद एक सबसिस्टम है, जिसे एक फ्लैग द्वारा सक्षम किया जाता है।
docker run -d --name nats \
-p 4222:4222 -p 127.0.0.1:8222:8222 \
-v nats-data:/data \
nats:2.14 -js -sd /data -m 8222-sd /data स्टोर डायरेक्टरी को सेट करता है। यदि आप इसे छोड़ देते हैं, तो JetStream अपना डेटा /tmp के अंतर्गत स्टोर करता है, जो सुनने में जितना टिकाऊ लगता है, उतना ही है। nats-box इमेज में आने वाले CLI का उपयोग करके एक स्ट्रीम बनाएँ।
docker run --rm -it --network host natsio/nats-box:latest \
nats stream add ORDERS \
--subjects 'orders.>' \
--storage file \
--retention limits \
--max-age 72h \
--max-bytes=1073741824 \
--discard old \
--defaultsवहाँ मौजूद प्रत्येक सीमा एक छोटे सर्वर पर अपना महत्व रखती है। --storage file वह है जो क्रैश के बाद भी सुरक्षित रहता है, क्योंकि मेमोरी स्ट्रीम ऐसा नहीं करती है। --max-bytes=1073741824 स्ट्रीम को 1 GiB पर सीमित करता है जिसे बाइट काउंट के रूप में लिखा जाता है, और --discard old सीमा तक पहुँचने पर नए राइट्स को अस्वीकार करने के बजाय सबसे पुराने संदेशों को हटा देता है। यदि आप इस सीमा को नहीं लगाते हैं, तो एक अनियंत्रित पब्लिशर डिस्क को भर देगा, जिस बिंदु पर आपका डेटाबेस भी रुक जाएगा, क्योंकि वे उसी डिस्क को साझा करते हैं।
एक durable consumer स्ट्रीम में अपनी स्थिति बनाए रखता है और रीस्टार्ट के बाद भी उसे याद रखता है। कंज्यूमर पर --max-deliver सेट करें ताकि जो संदेश हमेशा विफल रहता है, उसे बार-बार डिलीवर न किया जाए। जब किसी संदेश की डिलीवरी समाप्त हो जाती है, तो JetStream $JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.> पर एक एडवाइजरी पब्लिश करता है, और उस सब्जेक्ट को सब्सक्राइब करना ही वह तरीका है जिससे आप डेड लेटर पाथ बनाते हैं, जो RabbitMQ में एक फीचर के रूप में मिलता है। यह वह वास्तविक कार्य है जिसे आप स्वयं लिखते हैं।
बैकलॉग देखने के लिए, स्टोर किए गए संदेशों की संख्या के लिए nats stream report चलाएँ और प्रति कंज्यूमर बकाया एकनॉलेजमेंट और अनप्रोसेस्ड संदेशों के लिए nats consumer report ORDERS चलाएँ। अनप्रोसेस्ड संख्या वह है जिस पर आपको अलर्ट सेट करना चाहिए। डिस्क की लागत du -sh के साथ स्टोर डायरेक्टरी पर देखी जा सकती है, और यह तब तक बढ़ती है जब तक कि रिटेंशन लिमिट इसे कम न कर दे।
RabbitMQ: प्रत्येक message को acknowledge करें, विफलताओं को park करें
docker run -d --name rabbitmq \
-p 5672:5672 -p 127.0.0.1:15672:15672 \
-v rabbitmq-data:/var/lib/rabbitmq \
rabbitmq:4-managementअगस्त 2026 तक, वर्तमान series 4.3 है। Port 5672 AMQP (advanced message queuing protocol) के लिए है और 15672 management interface के लिए है। Interface को localhost पर रखें और SSH tunnel के माध्यम से उस तक पहुँचें।
Queue को x-queue-type argument को quorum पर सेट करके declare करें; default अभी भी classic है। Quorum queues हमेशा durable होती हैं और कुछ भी करने से पहले डेटा को disk पर लिखती हैं, इसलिए एक node पर आपको durable और transient विकल्पों के matrix के बजाय एक स्पष्ट व्यवहार मिलता है। Dead letter target को एक policy के साथ सेट करें।
docker exec rabbitmq rabbitmqctl set_policy DLX ".*" \
'{"dead-letter-exchange":"my-dlx", "dead-letter-routing-key":"my-routing-key"}' \
--apply-to queues --priority 7एक message चार कारणों से dead letter होता है: consumer इसे basic.reject या basic.nack के साथ reject करता है और requeue को false पर सेट करता है, इसका per-message TTL (time to live) समाप्त हो जाता है, queue लंबाई की सीमा पार कर जाती है, या यह quorum queue delivery limit से अधिक हो जाता है। RabbitMQ 4.0 के बाद से वह सीमा default रूप से 20 है, इसलिए जो handler error देता है और nack करता है, वह बीस बार retry करता है और फिर message को loop में डालने के बजाय dead letter exchange को सौंप देता है।
Memory वह जगह है जहाँ RabbitMQ छोटे VPS पर लोगों को हैरान कर देता है। Default high watermark उपलब्ध RAM का 0.6 है, और जब node इसे पार करता है, तो RabbitMQ publish करने वाले प्रत्येक connection को block कर देता है। आपके application को कोई error प्राप्त नहीं होता है। इसे एक ऐसा publish प्राप्त होता है जो कभी वापस नहीं आता, जिसे आपके अपने code में hang के रूप में देखा जाता है। Startup log उस संख्या को print करता है जिसकी गणना node ने की है:
Memory high watermark set to 1024 MiB (1073741824 bytes) of 8192 MiB (8589934592 bytes) totalDisk alarm भी इसी तरह publishers को block करता है जब free space default रूप से 50 MB से नीचे गिर जाता है। Quorum queues इसके ऊपर अपना गणित जोड़ती हैं: documentation प्रति message कम से कम 32 bytes in-memory metadata का बजट रखती है, जो प्रति 30,000 messages पर लगभग 1 MB है, और RAM में कम से कम तीन गुना प्रभावी write-ahead log size की सिफारिश करती है। WAL limit default रूप से 512 MiB है, इसलिए केवल वह सिफारिश ही 1.5 GB की मांग करती है। 2 GB के सर्वर पर, इसे rabbitmq.conf में कम करें, बजाय इसके कि यह उम्मीद करें कि default फिट हो जाएगा।
raft.wal_max_size_bytes = 64000000
vm_memory_high_watermark.relative = 0.5Backlog दो संख्याओं का योग है, और यह जोड़ी आपको बताती है कि आपके पास कौन सी विफलता है।
docker exec rabbitmq rabbitmqctl list_queues name messages messages_ready messages_unacknowledgedmessages_ready एक consumer की प्रतीक्षा कर रहा है। messages_unacknowledged deliver किया गया था और कभी ack नहीं किया गया। Flat ready count के बगल में बढ़ती हुई unacknowledged संख्या का मतलब है कि आपके workers ने jobs ले ली हैं और उन्हें पूरा करना बंद कर दिया है, जो कि केवल पीछे चल रही queue से एक अलग bug है।
एक ही बॉक्स पर Kafka, और यह कब तर्कसंगत नहीं रहता
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
bin/kafka-storage.sh format --standalone -t $KAFKA_CLUSTER_ID -c config/server.properties
bin/kafka-server-start.sh config/server.propertiesयह Kafka 4.3.1 के लिए क्विकस्टार्ट है, जो अगस्त 2026 तक अद्यतित है और KRaft मोड (Kafka Raft, जो Kafka 4.0 में ZooKeeper की जगह लेने वाला इन-बिल्ट कंट्रोलर है) में चल रहा है। इसका कंटेनर समकक्ष apache/kafka:4.3.1 है।
यदि आप स्वयं सेट नहीं करते हैं, तो स्टार्ट स्क्रिप्ट export KAFKA_HEAP_OPTS="-Xmx1G -Xms1G" को सेट कर देती है। इसलिए, ब्रोकर एक भी मैसेज स्टोर करने से पहले 1 GB Java heap रिज़र्व कर लेता है, और यह उस पेज कैश के लिए अतिरिक्त खाली RAM की अपेक्षा करता है जिससे यह डेटा पढ़ता है। 2 GB के VPS पर आपका एप्लिकेशन बची हुई RAM के लिए JVM के साथ प्रतिस्पर्धा करता है।
रिटेंशन (Retention) अगला आश्चर्य है। log.retention.hours का डिफ़ॉल्ट मान 168 है, जो सात दिन होता है, और log.retention.bytes का डिफ़ॉल्ट मान -1 है, जिसका अर्थ है कि आकार की कोई सीमा नहीं है। Kafka मैसेज को पूरे समय के लिए रखता है, चाहे हर कंज्यूमर ने उन्हें पढ़ा हो या नहीं। यही वह फीचर है जिसके लिए आप आए थे, और एक छोटी डिस्क पर यही विफलता का कारण भी बन सकता है, इसलिए इससे पहले कि आपको पता चले, प्रति टॉपिक बाइट सीमा सेट कर लें।
अब ईमानदारी की बात। एक सिंगल ब्रोकर का मतलब है रेप्लिकेशन फैक्टर 1, इसलिए acks=all एक डिस्क पर एक fsync में बदल जाता है। आपको एक मशीन की ड्यूरेबिलिटी मिलती है, साथ ही JVM ब्रोकर और कंट्रोलर की परिचालन लागत भी। पार्टिशन आपको उन ब्रोकर्स के बीच पैरेललिज्म (parallelism) देते हैं जो आपके पास हैं ही नहीं। रेप्लिकेशन, रैक अवेयरनेस और बाकी फ्लीट फीचर्स निष्क्रिय रहते हैं। JetStream आपको उसी बॉक्स पर कम मेमोरी में वही ड्यूरेबल रिप्ले देता है। यहाँ Kafka को सही ठहराने के दो ही कारण हैं: या तो कोई डाउनस्ट्रीम टूल केवल Kafka प्रोटोकॉल का उपयोग करता है (Debezium के साथ change data capture, या एनालिटिक्स लोडर), या आप प्रोडक्शन टोपोलॉजी को छोटे स्तर पर दोहरा रहे हैं। क्लस्टर में बढ़ने की योजना का मतलब है और अधिक मशीनें खरीदना, और तब तक यह सौदा वैसा ही है जैसा एक सिंगल नोड पर k3s चलाना, जहाँ आप एक नोड की विश्वसनीयता के लिए क्लस्टर की जटिलता का भुगतान करते हैं।
Kafka में बैकलॉग का मतलब कंज्यूमर लैग (consumer lag) है।
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group my-groupLAG कॉलम पढ़ें, जो प्रत्येक पार्टिशन के लिए LOG-END-OFFSET माइनस CURRENT-OFFSET है। यदि एक पार्टिशन पर लैग बढ़ रहा है जबकि अन्य स्थिर हैं, तो यह एक असमान की (uneven key) की ओर इशारा करता है, क्योंकि समान की वाले सभी मैसेज एक ही पार्टिशन पर जाते हैं और एक कंज्यूमर उन्हें अकेले हैंडल करता है।
जब बॉक्स रीस्टार्ट होता है तो क्या होता है
Core NATS उड़ान के दौरान मौजूद सभी डेटा खो देता है और तुरंत वापस आ जाता है, क्योंकि इसमें रिकवर करने के लिए कुछ भी नहीं होता है। JetStream स्टोर डायरेक्टरी से स्ट्रीम और कंज्यूमर पोजीशन को रीलोड करता है, इसलिए कंज्यूमर उसी ऑफसेट से फिर से शुरू हो जाते हैं जहाँ वे थे। RabbitMQ डिस्क से कोरम क्यू (quorum queues) को रिकवर करता है, जबकि क्लासिक ट्रांजिएंट क्यू और बिना पर्सिस्टेंट डिलीवरी मोड के पब्लिश किए गए कोई भी मैसेज समाप्त हो जाते हैं। Kafka स्टार्टअप पर अपने लॉग सेगमेंट को रीप्ले करता है, और अनक्लीन शटडाउन के बाद, ब्रोकर द्वारा कनेक्शन स्वीकार करने से पहले उस रिकवरी स्कैन में छोटी डिस्क पर भी कुछ मिनट लग सकते हैं।
दो चीजें एक बार सेट करना उचित है। कंटेनर को रीस्टार्ट पॉलिसी (restart: unless-stopped) दें या systemd यूनिट को इनेबल करें, ताकि कर्नल अपग्रेड रीबूट के बाद ब्रोकर आपके बिना वापस आ जाए। फिर क्रम को संभालें: एक ब्रोकर जो आपके एप्लिकेशन के बीस सेकंड बाद रेडी होता है, वह पहले कनेक्शन को रिफ्यूज कर देगा, और कुछ क्लाइंट लाइब्रेरी रीट्राय करने के बजाय एग्जिट हो जाती हैं। Compose healthchecks that hold a dependent service back until the broker is ready का उपयोग करके ब्रोकर पर ऐप को गेट करें।
अपने VPS पर लागत, अनुमानित के बजाय मापी गई
प्रकाशित थ्रूपुट आंकड़े ऐसे हार्डवेयर पर मापे जाते हैं जो आपके पास नहीं है, आमतौर पर स्थानीय NVMe वाले मल्टी-कोर सर्वर पर। उन्हें एक ऊपरी सीमा मानें और अपने बॉक्स पर स्वयं मापें।
docker stats --no-stream
free -m
sudo du -sh /var/lib/docker/volumes/*/_dataइन्हें तब चलाएं जब ब्रोकर निष्क्रिय (idle) हो, और फिर अपने वास्तविक ट्रैफिक के तहत दोबारा चलाएं। दोनों के बीच का अंतर वह संख्या है जो यह तय करती है कि ब्रोकर आपके एप्लिकेशन के साथ फिट बैठता है या नहीं। थ्रूपुट के एक मोटे अनुमान के लिए, किसी के ब्लॉग पोस्ट के बजाय प्रत्येक प्रोजेक्ट के अपने लोड जनरेटर का उपयोग करें: NATS के लिए nats bench pub test --msgs 100000 --clients 2, Kafka के लिए bin/kafka-producer-perf-test.sh, और RabbitMQ के लिए PerfTest। जनरेटर को उसी VPS पर चलाने से ब्रोकर और जनरेटर दोनों एक साथ माप लिए जाते हैं, जो ठीक है बशर्ते आप संख्या रिपोर्ट करते समय इसका उल्लेख करें।
एक सीमा सभी पर लागू होती है। यहाँ प्रत्येक टिकाऊ (durable) विकल्प fsync पर प्रतीक्षा करता है, इसलिए नेटवर्क-अटैच्ड स्टोरेज वाले VPS पर डिस्क ही सीमा निर्धारित करती है, और ब्रोकर बदलने से इसमें कोई सुधार नहीं होगा।
तीन वर्कलोड और प्रत्येक के लिए आवश्यक मैसेज क्यू
- एक वेब एप्लिकेशन के लिए बैकग्राउंड जॉब्स, जैसे ईमेल भेजना, इमेज का आकार बदलना, या वेबहुक डिलीवर करना। Postgres और
SKIP LOCKEDके साथ शुरुआत करें। जब आपको प्रति-मैसेज acks, डिलीवरी लिमिट और एक डेड लेटर क्यू की आवश्यकता हो जिसे आप स्वयं लॉजिक लिखे बिना देख सकें, या जब जॉब टेबल डेटाबेस की सबसे व्यस्त टेबल बन जाए, तो क्वोरम क्यू (quorum queues) के साथ RabbitMQ पर स्विच करें। - ऐसी घटनाएं जिन पर कई आंतरिक सेवाएं प्रतिक्रिया देती हैं, जहाँ खोए हुए मैसेज को जल्दी से नए मैसेज से बदल दिया जाता है। Core NATS का उपयोग करें, जहाँ routing scheme के रूप में subjects और जहाँ काम साझा करने की आवश्यकता हो वहां queue groups का उपयोग करें। उन विषयों (subjects) के सीमित सेट के लिए JetStream stream जोड़ें जिन्हें रीस्टार्ट के बाद भी सुरक्षित रहना आवश्यक है, और बाकी को मेमोरी में रहने दें।
- एक इवेंट लॉग जिसे कंज्यूमर शुरुआत से पढ़ते हैं, ऑडिट ट्रेल के लिए, रीड मॉडल को फिर से बनाने के लिए, या बाद में एनालिटिक्स को फीड करने के लिए। फाइल स्टोरेज और स्पष्ट बाइट कैप (byte cap) के साथ JetStream का उपयोग करें। Kafka को केवल तब चुनें जब कोई डाउनस्ट्रीम टूल Kafka प्रोटोकॉल की मांग करे, और उस अनुकूलता की कीमत के रूप में JVM heap को स्वीकार करें।
एक सर्वर पर गलत चुनाव की कीमत थ्रूपुट (throughput) नहीं है। इसकी कीमत सुबह तीन बजे होने वाली रिकवरी है, जब आपको यह जानने की आवश्यकता होती है कि क्या मैसेज अभी भी मौजूद हैं। उसी आधार पर चुनाव करें।
FAQ
क्या मैं 2 GB VPS पर Kafka चला सकता हूँ?
यह start तो हो जाएगा, लेकिन संसाधन बहुत सीमित रहेंगे। bin/kafka-server-start.sh तब KAFKA_HEAP_OPTS="-Xmx1G -Xms1G" सेट करता है जब आपने इसे override नहीं किया होता है, इसलिए कोई भी message स्टोर करने से पहले ही JVM 1 GB memory ले लेता है, जबकि Kafka को page cache के लिए इसके अतिरिक्त खाली memory की आवश्यकता होती है। यदि आप उसी सर्वर पर अपना application और database भी चलाते हैं, तो system swap का उपयोग करने लगेगा। साथ ही, आपको replication factor 1 ही मिलेगा, जिसका अर्थ है कि acks=all एक डिस्क पर एक fsync है, यानी आप Kafka की durability model के बिना ही उसकी operational लागत चुका रहे हैं। NATS JetStream उसी hardware पर बहुत कम memory में durable replay की सुविधा देता है।
क्या मुझे message queue की आवश्यकता है यदि मैं पहले से ही Postgres चला रहा हूँ?
अक्सर नहीं। एक job table read जिसमें transaction के अंदर SELECT ... FOR UPDATE SKIP LOCKED का उपयोग किया गया हो, वह at-least-once delivery, सुरक्षित concurrent workers, retries और dead letter table प्रदान करता है। इसमें आपको किसी अतिरिक्त service को monitor करने की आवश्यकता नहीं होती और आप वही backups उपयोग कर सकते हैं जो आप पहले से ले रहे हैं। इसे बदलने के संकेत स्पष्ट हैं: जब queue table आपका सबसे भारी write load बन जाए और autovacuum पीछे छूटने लगे, जब लंबे समय तक चलने वाले jobs transactions को खुला रखकर पूरे database के लिए vacuum को block कर दें, या जब किसी दूसरी service को स्वतंत्र रूप से उन्हीं events को consume करने की आवश्यकता हो।
क्या मुझे background jobs के लिए NATS JetStream का उपयोग करना चाहिए या RabbitMQ का?
यदि आप प्रति-message acknowledgement, delivery limit और dead letter routing जैसी सुविधाएँ built-in चाहते हैं, तो RabbitMQ चुनें। Quorum queues हमेशा durable होती हैं, RabbitMQ 4.0 से delivery limit डिफ़ॉल्ट रूप से 20 है, और एक policy के माध्यम से आप समाप्त हो चुके messages को dead letter exchange में भेज सकते हैं जिसे बाद में drain और inspect किया जा सकता है। यदि उन्हीं events को बाद में अन्य consumers द्वारा replay करने की आवश्यकता है, तो JetStream चुनें, क्योंकि stream acknowledgement के बाद भी messages को सुरक्षित रखता है जबकि queue ऐसा नहीं करती। JetStream के साथ आप --max-deliver सेट करते हैं और $JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.> advisory से dead letter path खुद तैयार करते हैं।
मुझे कैसे पता चलेगा कि मेरे consumers कितने पीछे हैं?
प्रत्येक broker के पास इसके लिए एक command है। RabbitMQ के लिए, rabbitmqctl list_queues name messages messages_ready messages_unacknowledged उस काम को अलग करता है जो consumer का इंतज़ार कर रहा है और उस काम को जो deliver हो चुका है पर ack नहीं हुआ है। JetStream के लिए, nats consumer report <stream> प्रति consumer unprocessed messages और outstanding acknowledgements दिखाता है। Kafka के लिए, kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group <group> प्रति partition एक LAG column प्रिंट करता है। Core NATS में पढ़ने के लिए कोई backlog नहीं होता, क्योंकि यह कुछ भी स्टोर नहीं करता है। इसके बजाय http://localhost:8222/varz पर slow_consumers counter को देखें: यह उन connections की गिनती करता है जिन्हें सर्वर ने पीछे छूट जाने के कारण बंद कर दिया है, जो कि message loss का संकेत है।