একটি VPS-এ NATS, RabbitMQ নাকি Kafka: কোনটি সেরা?
একটি সিঙ্গেল সার্ভারে মেসেজ কিউ নির্বাচনের সঠিক উপায় জানুন। ডেলিভারি গ্যারান্টি, মেমোরি খরচ এবং রিস্টার্ট বিহেভিয়ার বিশ্লেষণ করে দেখুন কেন Postgres অনেক ক্ষেত্রে সেরা সমাধান।
একটি সার্ভারের জন্য সংক্ষিপ্ত উত্তর
একটি VPS-এ মেসেজ কিউ ব্যবহারের সিদ্ধান্তটি মূলত ডেলিভারি গ্যারান্টির ওপর নির্ভর করে, গতির ওপর নয়। একটি সিঙ্গেল বক্সে ব্রোকার খুব কমই বটলনেক বা গতির প্রতিবন্ধক হয়, কারণ আপনার অ্যাপ্লিকেশন কোড, ডাটাবেস এবং ডিস্কের সীমাবদ্ধতা আগেই সামনে চলে আসে। এমন একটি টুল বেছে নিন যার ফেইলিয়র বিহেভিয়ার বা ব্যর্থতার ধরন আপনি মেনে নিতে পারবেন, এরপর আপনার বর্তমান সার্ভারের সক্ষমতা পরিমাপ করুন।
চারটি বিকল্প নিচে দেওয়া হলো, যা অধিকাংশ পাঠকের বিবেচনা করা উচিত:
- আপনার বর্তমান ডাটাবেসটিই ব্যবহার করুন। Postgres এবং
SELECT ... FOR UPDATE SKIP LOCKEDএকটি কার্যকর জব কিউ হিসেবে কাজ করে এবং এটি মনিটর করার জন্য নতুন কোনো প্রসেস যোগ করতে হয় না। - RabbitMQ ব্যবহার করুন যখন প্রতিটি মেসেজ এমন একটি কাজের একক যা অবশ্যই অ্যাকনলেজ (acknowledged) হতে হবে, নির্দিষ্ট সংখ্যক বার রিট্রাই (retried) করতে হবে এবং পরবর্তীতে এমন কোথাও জমা রাখতে হবে যেখানে মানুষ তা পরীক্ষা করতে পারে।
- NATS ব্যবহার করুন যখন মেসেজগুলো এমন ইভেন্ট হিসেবে কাজ করে যা আপনার সিস্টেমের বিভিন্ন অংশকে সক্রিয় করে। যেসব ইভেন্ট রিবুটের পরেও টিকে থাকা প্রয়োজন, সেগুলোর জন্য JetStream চালু করুন।
- Kafka ব্যবহার করুন যখন কোনো ডাউনস্ট্রিম টুল শুধুমাত্র Kafka প্রোটোকল সমর্থন করে। একটি সার্ভারের ক্ষেত্রে এটিই একমাত্র গ্রহণযোগ্য কারণ হতে পারে।
এই গাইডের বাকি অংশে এর যৌক্তিকতা ব্যাখ্যা করা হয়েছে: একটি ছোট VPS-এ প্রতিটি অপশনের জন্য মেমোরি ও ডিস্কের খরচ কত, বক্স রিবুট হলে কী ঘটে এবং ব্যবহারকারীরা সমস্যা টের পাওয়ার আগেই ব্যাকলগ দেখার সঠিক কমান্ড কী।
ডেলিভারি গ্যারান্টি বলতে আসলে কী বোঝায়
At most once মানে হলো ব্রোকার মেসেজটি পাঠিয়ে দেয় এবং তা ভুলে যায়। যদি কোনো কনজিউমার সংযুক্ত না থাকে, অথবা কাজ শেষ হওয়ার আগেই কোনো কনজিউমার বন্ধ হয়ে যায়, তবে মেসেজটি হারিয়ে যায় এবং এর কোনো রিপোর্ট পাওয়া যায় না।
At least once মানে হলো কাজ সফলভাবে সম্পন্ন হওয়ার পর কনজিউমার একটি একনলেজমেন্ট (ack) পাঠায়। যতক্ষণ না সেই ack পৌঁছায়, ব্রোকার মেসেজটি নিজের কাছে রাখে এবং পুনরায় ডেলিভারি করার চেষ্টা করে। এই পুনরায় ডেলিভারির কারণেই আপনার হ্যান্ডলারগুলোকে অবশ্যই idempotent হতে হবে: একই মেসেজ দুইবার প্রসেস করার ফলে যেন কার্ড থেকে দুইবার টাকা না কাটে। এন্ড-টু-এন্ড Exactly once কোনো ব্রোকার আপনাকে সরাসরি দেয় না। এটি আপনার নিজস্ব ডাটাবেসে থাকা একটি ইউনিক কি (unique key) থেকে নিশ্চিত করতে হয়।
Replay একটি আলাদা বৈশিষ্ট্য। একটি কিউ (queue) থেকে মেসেজ একনলেজ হওয়ার সাথে সাথে তা মুছে ফেলা হয়। অন্যদিকে একটি লগ (log) মেসেজটিকে একটি নির্দিষ্ট রিটেনশন উইন্ডো পর্যন্ত সংরক্ষণ করে, যাতে নতুন কোনো কনজিউমার শুরু থেকে পুরো হিস্ট্রি পড়তে পারে। Kafka এবং NATS JetStream হলো লগ। RabbitMQ হলো একটি কিউ। থ্রুপুটের চেয়ে এই পার্থক্যটিই আর্কিটেকচারের ওপর বেশি প্রভাব ফেলে।
Dead lettering হলো সেই প্রক্রিয়া যা বারবার ব্যর্থ হওয়া মেসেজগুলোর ক্ষেত্রে ঘটে। এটি না থাকলে, একটি পয়জন মেসেজ (poison message) লুপে আটকে থাকে এবং এই লুপটিকে দেখে মনে হয় ওয়ার্কারটি ব্যস্ত, অথচ আসলে সেটি অকেজো হয়ে আছে।
Postgres দিয়ে শুরু করুন এবং ব্রোকারকে তার সক্ষমতা প্রমাণ করতে দিন
অধিকাংশ একক-অ্যাপ্লিকেশন ওয়ার্কলোড দিনে কয়েক হাজার ব্যাকগ্রাউন্ড জব নিয়ে গঠিত। এটি একটি টেবিলের মধ্যেই অনায়াসে রাখা সম্ভব।
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);একজন ওয়ার্কার একটি ট্রানজ্যাকশনের ভেতরে একটি জব দাবি করে।
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 হলো পুরো কৌশলের মূল ভিত্তি। এটি যে সারিটি রিটার্ন করে সেটিকে লক করে ফেলে এবং অন্য কোনো ট্রানজ্যাকশন ইতিমধ্যে লক করে রেখেছে এমন সারিগুলোকে এড়িয়ে যায়, ফলে দুইজন ওয়ার্কার কখনোই একই জব দাবি করতে পারে না। যদি কোনো ওয়ার্কার ক্র্যাশ করে, তবে Postgres তার ট্রানজ্যাকশন বাতিল (abort) করে দেয়, লকটি মুক্ত হয়ে যায় এবং সারিটি পরবর্তী ওয়ার্কারের জন্য দৃশ্যমান হয়। আপনি পাচ্ছেন at-least-once ডেলিভারি, attempts বাড়িয়ে রিট্রাই করার সুবিধা এবং একটি ডেড লেটার টেবিল—সবই সেই স্থায়িত্বের (durability) ওপর ভিত্তি করে যার জন্য আপনি ইতিমধ্যে মূল্য পরিশোধ করছেন। ব্যাকলগ দেখার কুয়েরিটি হলো: SELECT count(*) FROM job WHERE run_after <= now();
যেখানে এটি কাজ করা বন্ধ করে দেয়। প্রতিটি দাবি (claim) এবং ডিলিট একটি রাইট অপারেশন, তাই উচ্চ হারের জব থাকলে তা ডেড রো ভার্সন রেখে যায় এবং কিউ টেবিলটি এমন একটি ক্লাসিক উদাহরণ যেখানে bloat, autovacuum-এর গতির চেয়ে দ্রুত বেড়ে যায়। দীর্ঘস্থায়ী জবগুলো পরিস্থিতি আরও খারাপ করে, কারণ কাজের সময়জুড়ে খোলা থাকা ট্রানজ্যাকশন পুরো ডেটাবেসের জন্য vacuum horizon-কে আটকে রাখে। পোলিং ল্যাটেন্সি বাড়িয়ে দেয় এবং NOTIFY এর সাথে LISTEN পোলিং দূর করলেও রাইট অপারেশনগুলো কমায় না। যখন জব টেবিলটি আপনার সবচেয়ে ব্যস্ত টেবিল হয়ে ওঠে, অথবা যখন দ্বিতীয় কোনো সার্ভিসের একই ইভেন্টের প্রয়োজন হয়, তখন কাজটিকে সরিয়ে নিন। এই সিদ্ধান্তটি ডেটাবেস কীভাবে ডিপ্লয় করা হয়েছে তার ওপর নির্ভর করে, তাই ডেটাবেসটি Docker-এ নাকি host-এ চলবে তা ব্রোকার যোগ করার আগেই নিশ্চিত করুন।
Redis হলো অন্য একটি টুল যা আপনি হয়তো ইতিমধ্যে ব্যবহার করছেন। Redis Streams আপনাকে XADD এবং XREADGROUP সহ কনজিউমার গ্রুপ, প্রতি গ্রুপের জন্য একটি পেন্ডিং লিস্ট এবং মৃত কনজিউমারের কাছ থেকে কাজ ফিরিয়ে নেওয়ার জন্য XAUTOCLAIM প্রদান করে। এটি ছোট এবং দ্রুত। একটি বক্সে ব্যবহারের ক্ষেত্রে সত্য বিষয়টি হলো: সাধারণ appendfsync everysec সেটিংয়ের কারণে, বিদ্যুৎ বিভ্রাটে প্রায় এক সেকেন্ডের রাইট ডেটা হারিয়ে যেতে পারে। ক্যাশ ইনভ্যালিডেশনের জন্য এটি ঠিক থাকলেও পেমেন্টের মতো কাজের জন্য এটি ভুল। যদি আপনার অ্যাপ্লিকেশনটি VPS-এ প্রোডাকশনে SQLite ব্যবহারের ওপর ভিত্তি করে তৈরি একটি একক প্রসেস হয়, তবে একই claim-and-delete প্যাটার্ন কাজ করবে, যদিও SQLite-এ SKIP LOCKED এর কোনো সমতুল্য ব্যবস্থা নেই এবং প্রতিটি ওয়ার্কারকে একটি মাত্র রাইট লকের ওপর সিরিয়ালাইজ হতে হয়।
NATS core: মেমোরি ছাড়া সাবজেক্ট রাউটিং
docker run -d --name nats \
-p 4222:4222 -p 127.0.0.1:8222:8222 \
nats:2.14 -m 82222026 সালের আগস্ট মাস অনুযায়ী বর্তমান সার্ভার লাইন হলো 2.14। -m 8222 ফ্ল্যাগটি HTTP মনিটরিং পোর্ট চালু করে, যা ডিফল্টভাবে বন্ধ থাকে এবং এতে কোনো প্রমাণীকরণ (authentication) নেই, তাই উপরের নিয়ম অনুযায়ী এটিকে localhost-এ bind করুন।
Core NATS হলো at-most-once ডেলিভারি সিস্টেম এবং এটি কোনো কিছু সংরক্ষণ করে না। একজন প্রকাশক (publisher) কোনো সাবজেক্টে (যেমন: orders.created) বার্তা পাঠালে, যার ফিল্টার মিলে যায় এমন প্রতিটি গ্রাহক (subscriber) একটি করে কপি পায়। যদি কেউ সাবস্ক্রাইব করা না থাকে, তবে বার্তাটি মুছে ফেলা হয় এবং প্রকাশক কোনো ত্রুটি দেখতে পায় না, কারণ সার্ভার বাইটগুলো গ্রহণ করার সাথে সাথেই প্রকাশকের কাজ শেষ হয়ে যায়। একটি কিউ গ্রুপ (যেখানে একাধিক গ্রাহক একই গ্রুপ নাম শেয়ার করে) ব্যবহার করলে সার্ভার প্রতি বার্তার জন্য একজন সদস্যকে বেছে নেয়, যা কোনো কিউ সংরক্ষণ না করেই কাজের ভার ভাগ করে দেয়।
এর ফুটপ্রিন্ট হলো সাবস্ক্রিপশন স্টেট এবং প্রতিটি কানেকশনের জন্য একটি রাইট বাফার, তাই এটি বার্তার পরিমাণের চেয়ে কানেকশনের সংখ্যার ওপর ভিত্তি করে চলে এবং ডিস্কে কোনো কিছু জমা হয় না। রিস্টার্টের আচরণও এর ওপর নির্ভর করে: চলমান বার্তাগুলো হারিয়ে যায়, ক্লায়েন্টরা নিজেরাই পুনরায় কানেক্ট করে এবং পুনরুদ্ধারের জন্য কোনো ধাপের অপেক্ষা করতে হয় না।
এখানে পর্যবেক্ষণ করার মতো কোনো ব্যাকলগ নেই, তাই ডেটা হারানোর দিকে নজর রাখুন। যখন কোনো গ্রাহক সার্ভারের লেখার গতির চেয়ে ধীরগতিতে সকেট থেকে ডেটা পড়ে, তখন সেই ক্লায়েন্টের জন্য সার্ভারের বাফার পূর্ণ হয়ে যায়। যদি ক্লায়েন্ট রাইট ডেডলাইনের মধ্যে ডেটা পড়ে শেষ করতে না পারে, তবে সার্ভার পুরো কানেকশনটি বন্ধ করে দেয় এবং একটি কাউন্টার বৃদ্ধি করে।
curl -s http://localhost:8222/varz | jq '.slow_consumers, .connections, .in_msgs, .out_msgs'slow_consumers-এর মান ক্রমাগত বাড়তে থাকলে বুঝতে হবে বার্তাগুলো ড্রপ হচ্ছে, তাই এটি একবার চেক না করে এর ওপর অ্যালার্ট সেট করুন। Core NATS এমন বার্তার জন্য উপযুক্ত যার মান দ্রুত শেষ হয়ে যায়: যেমন কোনো মেট্রিক, উপস্থিতি আপডেট (presence update), অথবা ক্যাশ ইনভ্যালিডেশন যা পরবর্তী ইভেন্টের মাধ্যমে এমনিতেই প্রতিস্থাপিত হবে।
NATS JetStream: একই প্রসেসে ডিউরেবল স্ট্রিম এবং রিপ্লে
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 নতুন রাইট প্রত্যাখ্যান না করে ক্যাপ পূর্ণ হলে পুরনো মেসেজগুলো মুছে ফেলে। এই ক্যাপ বাদ দিলে একটি অনিয়ন্ত্রিত পাবলিশার পুরো ডিস্ক পূর্ণ করে ফেলতে পারে, যার ফলে আপনার ডেটাবেসও বন্ধ হয়ে যাবে, কারণ তারা একই ডিস্ক শেয়ার করে।
একটি ডিউরেবল কনজিউমার স্ট্রিমে তার নিজস্ব অবস্থান ধরে রাখে এবং রিস্টার্টের পরেও তা বজায় থাকে। কনজিউমারে --max-deliver সেট করুন যাতে কোনো মেসেজ বারবার ব্যর্থ হলে তা আর ডেলিভারি না হয়। যখন একটি মেসেজের ডেলিভারি সীমা শেষ হয়ে যায়, JetStream $JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>-এ একটি অ্যাডভাইজরি পাবলিশ করে। সেই সাবজেক্টে সাবস্ক্রাইব করার মাধ্যমেই আপনি ডেড লেটার পাথ তৈরি করতে পারেন, যা RabbitMQ-তে একটি ফিচার হিসেবে থাকে। এটি এমন একটি কাজ যা আপনাকে নিজে কোড করে তৈরি করতে হবে।
ব্যাকলগ দেখার জন্য, সংরক্ষিত মেসেজের সংখ্যার জন্য nats stream report এবং প্রতিটি কনজিউমারের বকেয়া অ্যাকনলেজমেন্ট ও প্রসেস না হওয়া মেসেজের জন্য nats consumer report ORDERS চালান। প্রসেস না হওয়া মেসেজের সংখ্যার ওপর ভিত্তি করেই অ্যালার্ম সেট করা উচিত। স্টোর ডিরেক্টরিতে du -sh চালিয়ে ডিস্কের ব্যবহার দেখা যায়, এবং রিটেনশন লিমিট দ্বারা ট্রিম না হওয়া পর্যন্ত এটি বাড়তেই থাকে।
RabbitMQ: প্রতিটি মেসেজ একনলেজ করুন, ব্যর্থ মেসেজগুলো আলাদা করুন
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 অনুযায়ী বর্তমান সিরিজটি হলো 4.3। পোর্ট 5672 হলো AMQP (advanced message queuing protocol) এবং 15672 হলো ম্যানেজমেন্ট ইন্টারফেস। ইন্টারফেসটিকে localhost-এ রাখুন এবং SSH tunnel-এর মাধ্যমে এটি ব্যবহার করুন।
x-queue-type আর্গুমেন্টটিকে quorum-এ সেট করে কিউ (queue) ঘোষণা করুন; ডিফল্ট মান এখনো classic। কোরাম কিউ (quorum queues) সবসময় টেকসই (durable) হয় এবং অন্য কোনো কাজ করার আগেই ডিস্কে ডেটা লিখে ফেলে। ফলে একটি নোডে আপনি টেকসই এবং অস্থায়ী অপশনের জটিলতার পরিবর্তে একটি স্পষ্ট আচরণ পাবেন। একটি পলিসির মাধ্যমে ডেড লেটার টার্গেট সেট করুন।
docker exec rabbitmq rabbitmqctl set_policy DLX ".*" \
'{"dead-letter-exchange":"my-dlx", "dead-letter-routing-key":"my-routing-key"}' \
--apply-to queues --priority 7একটি মেসেজ চারটি কারণে ডেড লেটার হিসেবে গণ্য হয়: কোনো কনজিউমার যদি basic.reject বা basic.nack ব্যবহার করে মেসেজটি রিজেক্ট করে এবং requeue-কে false সেট করে, মেসেজের TTL (time to live) শেষ হয়ে গেলে, কিউ তার দৈর্ঘ্যের সীমা অতিক্রম করলে, অথবা এটি কোরাম কিউ ডেলিভারি লিমিট অতিক্রম করলে। RabbitMQ 4.0 থেকে এই লিমিট ডিফল্টভাবে 20 নির্ধারণ করা হয়েছে। ফলে কোনো হ্যান্ডলার যদি বারবার ত্রুটি দেখায় এবং nack করে, তবে সেটি বিশবার পুনরায় চেষ্টা করার পর মেসেজটিকে লুপে না ফেলে ডেড লেটার এক্সচেঞ্জে পাঠিয়ে দেয়।
ছোট VPS-এ RabbitMQ ব্যবহারের ক্ষেত্রে মেমোরি একটি চমকপ্রদ বিষয়। ডিফল্ট হাই ওয়াটারমার্ক হলো উপলব্ধ RAM-এর 0.6 অংশ। নোডটি যখন এই সীমা অতিক্রম করে, RabbitMQ পাবলিশ করা প্রতিটি কানেকশন ব্লক করে দেয়। আপনার অ্যাপ্লিকেশন কোনো এরর মেসেজ পায় না। এটি এমন একটি পাবলিশ রিকোয়েস্ট পায় যা কখনোই শেষ হয় না, যা আপনার কোডে হ্যাং হিসেবে ধরা দেয়। স্টার্টআপ লগে নোডটি যে সংখ্যাটি গণনা করে তা প্রিন্ট হয়:
Memory high watermark set to 1024 MiB (1073741824 bytes) of 8192 MiB (8589934592 bytes) totalডিস্ক অ্যালার্ম একইভাবে পাবলিশারদের ব্লক করে যখন ফ্রি স্পেস ডিফল্টভাবে 50 MB-এর নিচে নেমে যায়। কোরাম কিউ এর সাথে নিজস্ব গাণিতিক হিসাব যোগ করে: ডকুমেন্টেশন অনুযায়ী প্রতি মেসেজের জন্য কমপক্ষে 32 বাইট ইন-মেমোরি মেটাডেটা প্রয়োজন, যা প্রতি 30,000 মেসেজের জন্য প্রায় 1 MB। এছাড়া কার্যকর write-ahead log সাইজের কমপক্ষে তিনগুণ RAM রাখার পরামর্শ দেওয়া হয়। WAL লিমিট ডিফল্টভাবে 512 MiB, তাই শুধুমাত্র এই পরামর্শ অনুযায়ীই 1.5 GB RAM প্রয়োজন। 2 GB সার্ভারে, ডিফল্ট সেটিংস কাজ করবে এমন আশা না করে rabbitmq.conf-এ মান কমিয়ে নিন।
raft.wal_max_size_bytes = 64000000
vm_memory_high_watermark.relative = 0.5ব্যাকলগ দুটি সংখ্যার সমষ্টি, এবং এই জোড়াটি আপনাকে বলে দেয় আপনি কোন ধরনের ব্যর্থতার সম্মুখীন হচ্ছেন।
docker exec rabbitmq rabbitmqctl list_queues name messages messages_ready messages_unacknowledgedmessages_ready হলো কনজিউমারের জন্য অপেক্ষমাণ মেসেজ। messages_unacknowledged হলো সেই মেসেজ যা ডেলিভারি হয়েছে কিন্তু কখনোই একনলেজ (ack) করা হয়নি। রেডি কাউন্ট স্থির থাকা অবস্থায় আন-একনলেজড কাউন্ট বাড়তে থাকলে বুঝতে হবে আপনার ওয়ার্কাররা কাজগুলো নিয়েছে কিন্তু শেষ করছে না। এটি এমন কিউ থেকে ভিন্ন সমস্যা যা কেবল পিছিয়ে পড়েছে।
একটি বক্সে 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 রিজার্ভ করে ফেলে এবং এটি পড়ার জন্য পেজ ক্যাশের (page cache) উদ্দেশ্যে অতিরিক্ত ফ্রি র্যাম প্রত্যাশা করে। একটি 2 GB VPS-এ আপনার অ্যাপ্লিকেশন তখন অবশিষ্ট র্যামের জন্য JVM-এর সাথে প্রতিযোগিতায় লিপ্ত হয়।
রিটেনশন (Retention) হলো পরবর্তী চমক। log.retention.hours-এর ডিফল্ট মান 168, অর্থাৎ সাত দিন, এবং log.retention.bytes-এর ডিফল্ট মান -1, যার অর্থ কোনো সাইজ লিমিট নেই। প্রতিটি কনজিউমার মেসেজগুলো পড়েছে কি না তা বিবেচনা না করেই Kafka পুরো সময়সীমার জন্য মেসেজগুলো রেখে দেয়। এটি সেই ফিচার যার জন্য আপনি এটি ব্যবহার করছেন, কিন্তু ছোট ডিস্কের ক্ষেত্রে এটিই ব্যর্থতার কারণ হতে পারে, তাই সমস্যা হওয়ার আগেই প্রতি টপিকের জন্য একটি বাইট লিমিট সেট করুন।
এখন আসল কথা। একটি সিঙ্গেল ব্রোকারের অর্থ হলো রেপ্লিকেশন ফ্যাক্টর 1, তাই acks=all একটি ডিস্কে একটি fsync-এর মাধ্যমে সম্পন্ন হয়। আপনি একটি মেশিনের স্থায়িত্ব পাচ্ছেন, কিন্তু এর বিপরীতে JVM ব্রোকার এবং একটি কন্ট্রোলারের পরিচালন ব্যয় বহন করতে হচ্ছে। পার্টিশনগুলো ব্রোকারগুলোর মধ্যে প্যারালালিজম সুবিধা দেয় যা আপনার নেই। রেপ্লিকেশন, র্যাক অ্যাওয়ারনেস এবং অন্যান্য ফ্লিট ফিচারগুলো এখানে নিষ্ক্রিয় থাকে। 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) রয়েছে, কারণ একই কি (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-এ খরচ, অনুমানের পরিবর্তে পরিমাপ করুন
প্রকাশিত থ্রুপুট (throughput) পরিসংখ্যান এমন হার্ডওয়্যারে পরিমাপ করা হয় যা আপনার কাছে নেই, সাধারণত লোকাল 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-এর জন্য অপেক্ষা করে, তাই নেটওয়ার্ক-সংযুক্ত স্টোরেজ (network-attached storage) বিশিষ্ট VPS-এ ডিস্কই সীমাবদ্ধতা নির্ধারণ করে, এবং ব্রোকার পরিবর্তন করলে এই সীমা বাড়বে না।
তিনটি ওয়ার্কলোড এবং প্রতিটি যে মেসেজ কিউ চায়
- একটি ওয়েব অ্যাপ্লিকেশনের জন্য ব্যাকগ্রাউন্ড জব, যেমন ইমেইল পাঠানো, ইমেজ রিসাইজ করা বা ওয়েবহুক ডেলিভারি করা। Postgres এবং
SKIP LOCKEDদিয়ে শুরু করুন। যখন আপনার প্রতি-মেসেজ একনলেজমেন্ট (acks), ডেলিভারি লিমিট এবং এমন একটি ডেড লেটার কিউ প্রয়োজন হয় যা নিজে লজিক না লিখেও পরীক্ষা করা যায়, অথবা যখন জব টেবিলটি ডাটাবেসের সবচেয়ে ব্যস্ত টেবিলে পরিণত হয়, তখন RabbitMQ এবং কোরাম কিউ (quorum queues)-তে চলে যান। - এমন ইভেন্ট যা একাধিক অভ্যন্তরীণ সার্ভিস গ্রহণ করে, যেখানে কোনো মেসেজ হারিয়ে গেলে দ্রুত নতুন একটি মেসেজ দিয়ে তা প্রতিস্থাপন করা যায়। এক্ষেত্রে Core NATS ব্যবহার করুন, যেখানে রাউটিং স্কিম হিসেবে সাবজেক্ট এবং কাজের ভাগাভাগির জন্য কিউ গ্রুপ ব্যবহার করা হয়। যে অল্প সংখ্যক সাবজেক্ট রিস্টার্টের পরেও টিকে থাকা প্রয়োজন, সেগুলোর জন্য JetStream স্ট্রিম যোগ করুন এবং বাকিগুলো মেমরিতে রাখুন।
- একটি ইভেন্ট লগ যা কনজিউমাররা শুরু থেকে পড়ে, অডিট ট্রেইল রাখা, রিড মডেল পুনর্গঠন বা পরবর্তীতে অ্যানালিটিক্সের জন্য। ফাইল স্টোরেজ এবং নির্দিষ্ট বাইট ক্যাপসহ JetStream ব্যবহার করুন। শুধুমাত্র তখনই Kafka বেছে নিন যখন কোনো ডাউনস্ট্রিম টুলের Kafka প্রোটোকল প্রয়োজন হয়, এবং সেই সামঞ্জস্যতার মূল্য হিসেবে JVM হিপ ব্যবহারের বিষয়টি মেনে নিন।
একটি সিঙ্গেল সার্ভারে ভুল নির্বাচনের মাশুল থ্রুপুট নয়। এর আসল মাশুল হলো রাত তিনটার সময় রিকভারি করা, যখন আপনাকে নিশ্চিত হতে হয় যে মেসেজগুলো আদৌ টিকে আছে কি না। সেই পরিস্থিতির কথা মাথায় রেখেই নির্বাচন করুন।
FAQ
আমি কি 2 GB VPS-এ Kafka চালাতে পারি?
এটি চালু হবে, তবে বেশ চাপের মুখে থাকবে। আপনি যদি ওভাররাইড না করেন তবে bin/kafka-server-start.sh ডিফল্টভাবে KAFKA_HEAP_OPTS="-Xmx1G -Xms1G" সেট করে, তাই কোনো মেসেজ জমা করার আগেই JVM 1 GB মেমোরি দখল করে নেয় এবং Kafka এর বাইরে পেজ ক্যাশের জন্য খালি মেমোরির ওপর নির্ভর করে। একই সার্ভারে আপনার অ্যাপ্লিকেশন এবং ডাটাবেস যোগ করলে তা সোয়াপ (swap) মেমোরি ব্যবহার করতে বাধ্য হবে। এছাড়া আপনি রেপ্লিকেশন ফ্যাক্টর 1 ব্যবহার করবেন, যার মানে acks=all হলো একটি ডিস্কে একটি fsync, অর্থাৎ আপনি Kafka-এর স্থায়িত্বের মডেল ছাড়াই এর পরিচালনার খরচ বহন করছেন। একই হার্ডওয়্যারে NATS JetStream অনেক কম মেমোরিতে টেকসই রিপ্লে (durable replay) সুবিধা দেয়।
আমার যদি ইতিমধ্যে Postgres থাকে, তবে কি আমার মেসেজ কিউ (message queue) প্রয়োজন?
সাধারণত প্রয়োজন নেই। একটি job টেবিল থেকে SELECT ... FOR UPDATE SKIP LOCKED ব্যবহার করে ট্রানজেকশনের মাধ্যমে রিড করলে আপনি 'এট-লিস্ট-ওয়ানস' (at-least-once) ডেলিভারি, নিরাপদ কনকারেন্ট ওয়ার্কার, রিট্রাই এবং ডেড লেটার টেবিল সুবিধা পাবেন, যার জন্য বাড়তি কোনো সার্ভিস মনিটর করতে হবে না এবং আপনার বর্তমান ব্যাকআপ পদ্ধতিই যথেষ্ট। কিউ থেকে বেরিয়ে আসার সংকেতগুলো হলো: কিউ টেবিলটি যদি আপনার সবচেয়ে ভারী রাইট লোড হয়ে দাঁড়ায় এবং অটোভ্যাকিউম (autovacuum) পিছিয়ে পড়ে, দীর্ঘস্থায়ী কাজগুলো যদি ট্রানজেকশন ওপেন রেখে পুরো ডাটাবেসের ভ্যাকুয়াম আটকে দেয়, অথবা যদি দ্বিতীয় কোনো সার্ভিসের স্বাধীনভাবে একই ইভেন্টগুলো কনজিউম করার প্রয়োজন হয়।
ব্যাকগ্রাউন্ড জবের জন্য আমার কি NATS JetStream নাকি RabbitMQ ব্যবহার করা উচিত?
যদি আপনি প্রতি-মেসেজ অ্যাকনলেজমেন্ট (per-message acknowledgement), ডেলিভারি লিমিট এবং বিল্ট-ইন ডেড লেটার রাউটিং চান, তবে RabbitMQ ব্যবহার করুন। কোরাম কিউ (Quorum queues) সবসময় টেকসই হয়, RabbitMQ 4.0 থেকে ডেলিভারি লিমিট ডিফল্টভাবে 20 সেট করা থাকে এবং একটি পলিসির মাধ্যমে ব্যর্থ মেসেজগুলোকে ডেড লেটার এক্সচেঞ্জে পাঠানো যায়, যা আপনি পরে পরীক্ষা করতে পারবেন। যদি একই ইভেন্টগুলো পরে অন্য কনজিউমারদের রিপ্লে করার প্রয়োজন হয়, তবে JetStream ব্যবহার করুন, কারণ কিউ মেসেজ ডিলিট করে দিলেও স্ট্রিম অ্যাকনলেজমেন্টের পরেও মেসেজগুলো রেখে দেয়। JetStream-এ আপনি --max-deliver সেট করতে পারেন এবং $JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.> অ্যাডভাইজরি থেকে ডেড লেটার পাথ নিজে তৈরি করতে পারেন।
আমার কনজিউমাররা কতটা পিছিয়ে আছে তা কীভাবে বুঝব?
প্রতিটি ব্রোকারের জন্য একটি নির্দিষ্ট কমান্ড রয়েছে। RabbitMQ-এর ক্ষেত্রে, rabbitmqctl list_queues name messages messages_ready messages_unacknowledged কনজিউমারের জন্য অপেক্ষমাণ কাজ এবং ডেলিভারি হওয়া কিন্তু অ্যাকনলেজ না হওয়া কাজের মধ্যে পার্থক্য দেখায়। JetStream-এর ক্ষেত্রে, nats consumer report <stream> প্রতি কনজিউমারে প্রসেস না হওয়া মেসেজ এবং বকেয়া অ্যাকনলেজমেন্ট দেখায়। Kafka-এর ক্ষেত্রে, kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group <group> প্রতি পার্টিশনে একটি LAG কলাম প্রিন্ট করে। কোর NATS-এ পড়ার মতো কোনো ব্যাকলগ নেই, কারণ এটি কিছুই সংরক্ষণ করে না; এর পরিবর্তে http://localhost:8222/varz-এ slow_consumers কাউন্টারটি দেখুন: এটি সেই সংযোগগুলোর সংখ্যা দেখায় যা সার্ভার পিছিয়ে পড়ার কারণে বন্ধ করে দিয়েছে, যার অর্থ হলো মেসেজ লস।