Apache Kafka পরিচিতি
এই পাঠে যা শিখবেন
- Kafka কেন একটি "log", "queue" নয় — এই পার্থক্যের গভীর তাৎপর্য
- Broker, cluster, partition, replica — চারটি core building block
- ZooKeeper থেকে KRaft-এ migration — কেন এই বদল
- Offset, retention, consumer group — কেন Kafka replay সমর্থন করে
১ · Kafka কী — তিন বাক্যে
Apache Kafka হলো একটি distributed event streaming platformDistributed Event Streamingএকাধিক machine-এ ছড়ানো একটি system যা continuously events publish/subscribe ও store করে। Real-time pipeline, microservices integration, ও event-driven architecture-এর ভিত্তি।। মূলত এটি একটি append-only distributed log — events শুধু শেষে যোগ হয়, কখনো পরিবর্তন বা মুছে ফেলা হয় না (retention period না পেরোনো পর্যন্ত)।
LinkedIn-এ Jay Kreps, Neha Narkhede, Jun Rao তৈরি করেন (২০১১), নাম "Franz Kafka"-র নামে। ২০১১-তে open source — আজ Fortune 500-এর ৮০%+ ব্যবহার করে।
Traditional queue (RabbitMQ): message একবার delivered → মুছে যায়। One-time consumption।
Kafka log: message লেখা থাকে দিন/সপ্তাহ। অনেক consumer একই message পৃথকভাবে পড়তে পারে। Replay সম্ভব — bug fix-এর পর গত ৭ দিনের data পুনরায় process।
২ · Why Kafka — কেন de-facto standard হলো
২০১৫-র পর Kafka real-time data-এর "universal substrate" হয়ে উঠেছে। কারণগুলো:
- Throughput: একটি broker ১ million+ messages/sec handle করে। ক্লাস্টারে ১০০ million+।
- Durability: Disk-এ persist, replicate। Broker fail করলেও data যায় না।
- Decoupling: Producer ও consumer একে অপরকে চেনে না। Producer down হলেও consumer পড়ে চলে।
- Scalability: Partition যোগ করে horizontal scale। Linear scaling।
- Replay: Consumer offset rewind করে আগের data পুনরায় পড়তে পারে।
- Ecosystem: Kafka Connect (১০০+ connector), Kafka Streams, ksqlDB — সম্পূর্ণ platform।
৩ · Architecture — broker, cluster, topic
Broker: একটি Kafka server। সাধারণত একটি machine = একটি broker। Broker messages store ও serve করে।
Cluster: একাধিক broker একসাথে — একটি Kafka cluster। Production-এ minimum ৩ broker (fault tolerance-এর জন্য)। বড় deployment-এ ১০০+ broker।
Topic: একটি logical category — events-এর "channel"। যেমন bkash-transactions, pathao-rides, daraz-orders। Producer এই topic-এ publish করে; consumer subscribe।
৪ · Partition — Kafka-এর scaling secret
একটি topic শুধু একটি broker-এ থাকলে — সেই broker-এর CPU/disk-ই সর্বোচ্চ throughput limit। তাই Kafka topic-কে partitionPartitionএকটি topic-কে অনেকগুলি ছোট log-এ ভাগ করার পদ্ধতি — প্রতিটি partition আলাদা broker-এ থাকতে পারে। এটাই Kafka-র parallelism ও scaling-এর মূল।-এ ভাগ করে — প্রতিটি partition একটি আলাদা log, যা ভিন্ন broker-এ থাকতে পারে।
প্রতিটি partition-এ message order রক্ষিত। কিন্তু partition-এর মধ্যে কোনো order guarantee নেই। তাই partition-key carefully বাছতে হয় — একই user-এর সব events একই partition-এ যাওয়া উচিত (causality preserve)।
$$\text{partition} = \text{hash}(\text{key}) \bmod \text{num\_partitions}$$
৫ · Replication — fault tolerance
প্রতিটি partition-এর replicas থাকে — সাধারণত ৩টি (replication factor = ৩)। একজন leader, বাকিরা followers। Producer leader-এ লেখে; followers async/sync replicate করে।
Leader broker fail করলে — followers থেকে নতুন leader নির্বাচিত হয় (automatic failover, সেকেন্ডে)। Data loss নেই (যদি acks=all)।
acks=all দেয়, message ISR-এর সব broker-এ পৌঁছালে তবে success। ১টি broker fail করলেও data safe।
৬ · ZooKeeper থেকে KRaft
Historically, Kafka cluster metadata (broker membership, topic config, leader election) ZooKeeperApache ZooKeeperএকটি distributed coordination service — distributed lock, configuration, leader election handle করে। Hadoop ecosystem-এর core component, কিন্তু operational overhead আছে।-এ store করত। আলাদা cluster, আলাদা ops burden।
KIP-500-এর মাধ্যমে (২০২০-২০২২) Kafka ZooKeeper বাদ দিল — KRaftKRaft (Kafka Raft)Kafka-র own consensus protocol (Raft algorithm-এ ভিত্তিক) — যা ZooKeeper-এর কাজ Kafka-র ভিতরেই করে। Kafka 3.3+ production-ready। mode-এ Kafka নিজেই Raft consensus algorithm দিয়ে metadata manage করে। Kafka 3.3+ (২০২২) production-ready।
KRaft-এর সুবিধা:
- Single system to operate — ZooKeeper deployment-এর ঝামেলা নেই।
- Faster controller failover (সেকেন্ডের চেয়ে দ্রুত)।
- Million-partition cluster scalable।
- Simpler security model।
৭ · Offset, retention, replay
প্রতিটি message-এ একটি sequential offset থাকে (0, 1, 2, …) প্রতিটি partition-এর মধ্যে। Consumer group তাদের নিজস্ব offset track করে — "আমি কোন position পর্যন্ত পড়েছি"।
Retention: default ৭ দিন। যেতে দিতে পারেন ৩০ দিন, ৯০ দিন, এমনকি infinite (retention.ms=-1)। Disk space-এর সাথে trade-off।
Replay-এর শক্তি: bug detect হলো — fraud detector গত ২ দিন miss করেছে। Kafka-তে consumer offset rewind করে গত ২ দিনের data পুনরায় process করুন। Database-based queue-এ এটি impossible।
৮ · Bangladesh fintech-এ Kafka adoption
আজকের Bangladesh-এর সব major fintech Kafka ব্যবহার করছে:
- bKash: transaction events, agent activity, fraud alerts — সব Kafka-তে।
- Nagad: real-time settlement, KYC verification pipeline।
- Pathao: ride events, driver location streams, food delivery।
- Daraz: order events, inventory sync, recommendation features।
- Banks (BRAC, City): CDC from core banking → fraud, AML, ML।
- BTRC: CDR aggregation থেকে DDoS detection।
৯ · Hands-on — Kafka local-এ চালানো
Docker দিয়ে single-node Kafka cluster (KRaft mode):
version: '3.8'
services:
kafka:
image: apache/kafka:3.7.0
container_name: kafka-kraft
ports:
- "9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
CLUSTER_ID: "MkU3OEVBNTcwNTJENDM2Qk"
docker-compose up -d চালালে — single-broker Kafka KRaft mode-এ চালু। ZooKeeper দরকার নেই। Production-এ অবশ্যই ৩+ broker।
Topic তৈরি ও producer/consumer test করুন CLI দিয়ে:
# Topic তৈরি — 3 partition, replication factor 1 (single broker)
docker exec -it kafka-kraft kafka-topics.sh \
--create --topic bkash-transactions \
--partitions 3 --replication-factor 1 \
--bootstrap-server localhost:9092
# Topic list দেখুন
docker exec -it kafka-kraft kafka-topics.sh \
--list --bootstrap-server localhost:9092
# Producer — terminal 1
docker exec -it kafka-kraft kafka-console-producer.sh \
--topic bkash-transactions --bootstrap-server localhost:9092
# এখন type করুন:
# {"user":"01711000001","amount":500,"type":"send_money"}
# Consumer — terminal 2
docker exec -it kafka-kraft kafka-console-consumer.sh \
--topic bkash-transactions --from-beginning \
--bootstrap-server localhost:9092
--from-beginning মানে topic-এর শুরু থেকে — Kafka-র replay capability প্রদর্শিত।
min.insync.replicas=2, acks=all। Monitoring (Prometheus + JMX exporter) ছাড়া Kafka চালাবেন না।
ভাবনার প্রশ্ন
প্রতিটি প্রশ্ন নিজে কিছুক্ষণ ভাবুন — তারপর "→ উত্তর" চাপুন।
প্র ০১ "Kafka একটি queue নয়, একটি log" — এই পার্থক্য কী impact ফেলে architecture decision-এ? RabbitMQ-র সাথে তুলনা করুন।
এই পার্থক্য সূক্ষ্ম কিন্তু architecture-এ revolutionary।
Traditional queue (RabbitMQ, ActiveMQ):
- Message deliver হলে — queue থেকে মুছে যায়। Storage-এ persist থাকে না বেশিক্ষণ।
- একটি message একজন consumer-ই পায় (or fanout)।
- Downstream system যোগ করতে — pipeline পরিবর্তন।
- Reprocess অসম্ভব — bug fix হলো, কিন্তু গত data পুনরায় process করা যাবে না।
Kafka log:
- Message লেখা থাকে retention period পর্যন্ত (দিন/সপ্তাহ/অসীম)।
- অনেক consumer group একই data পৃথকভাবে পড়তে পারে — সবাই নিজের offset।
- নতুন service add করতে শুধু নতুন consumer group চালু — Kafka-তে কিছু পরিবর্তন নেই।
- Replay possible — bug fix-এর পর সব consumer offset rewind, পুনরায় process।
Architecture impact:
(১) Event sourcing: Kafka-তে state-কে event log হিসেবে modeled করা যায়। "Account state at time T" = "সব events from t0 to T-এর projection"। Database-এর alternative — ACID সহকারে।
(২) Microservice integration: Service A → Kafka → Service B, C, D, E। Producer জানে না কে subscribe করছে। Loose coupling চরমে।
(৩) Time travel debugging: Production bug — Kafka offset rewind করে exactly যা ঘটেছিল replay। Database-এ এটি impossible।
(৪) Multi-tenant analytics: Same data → fraud detection (real-time), analytics (batch), ML training (অফলাইন), archive (S3) — সবাই একই Kafka topic পড়ছে।
(৫) Schema evolution: Schema Registry-এর সাথে — backward/forward compatible schemas। Old consumers নতুন data পড়তে পারে।
RabbitMQ এখনো relevant:
- Task queue (ছোট, ম্যানুয়াল priority)।
- Complex routing (header, topic exchange)।
- Per-message ack — granular control।
- Lower operational complexity ছোট scale-এ।
মূল উপলব্ধি: "Log" mental model architect-কে স্বাধীনতা দেয় — events একবার লিখলে চিরকাল available (retention period পর্যন্ত)। অনেক pattern (event sourcing, CDC, stream processing) এই ভিত্তিতে দাঁড়িয়ে। RabbitMQ-র জগৎ-এ এই pattern সম্ভব না।
প্র ০২ আপনি bKash-এ Kafka cluster design করছেন — দৈনিক ৪ কোটি transaction, ১০০ KB average message। কতগুলো broker, partition, replication factor, retention বাছবেন? রেট-লিমিট কোথায়?
এটি classic capacity planning question। সঠিক উত্তর সংখ্যাগত — অনুমান যাচাই করার সুযোগ।
সংখ্যা বিশ্লেষণ:
- ৪ কোটি/দিন = ৪×১০⁷ / ৮৬,৪০০ ≈ ৪৬৩ TPS গড়।
- Peak (lunch, evening): ৫× = ~২,৫০০ TPS।
- Message size: ১০০ KB → bandwidth: ২৫০ MB/sec peak।
- Daily volume: ৪×১০⁷ × ১০০ KB = ৪ TB।
Cluster sizing:
- Brokers: production minimum ৫ broker (rolling restart-এ ১ down হলেও RF=3 maintained)। AWS m6i.2xlarge প্রতিটি (৮ vCPU, ৩২ GB RAM, NVMe SSD)।
- Per-broker capacity: ১ million msg/sec, কিন্তু production conservative — ১০০K msg/sec budget।
- Replication factor: ৩। দু'টি AZ লিভিং, ২টি broker simultaneously fail-এও data safe।
- min.insync.replicas: ২।
acks=all-এ ১ replica down হলেও write হবে।
Topic ও partition design:
- Topic-এর লিস্ট:
transactions,user-events,agent-events,fraud-decisions,settlement-events। - Transactions topic partitions: ২৪। Per-partition target ১০০ msg/sec → ২৪×১০০=২,৪০০ TPS capacity (peak handle)।
- Partition key:
user_phone_number। একই user-এর সব transactions একই partition → causality। - Skew check: "VIP" user-এর partition hot হতে পারে — composite key (
phone + transaction_type) consider।
Retention strategy:
- transactions: ৭ দিন (replay বাজেট)। ৪ TB × ৭ × ৩ replicas = ৮৪ TB storage। NVMe-এ ব্যয়বহুল — tiered storage (Confluent) বা compaction।
- fraud-decisions: ৩ দিন।
- analytics-aggregations: ৩০ দিন।
- S3 archive nightly — Kafka Connect S3 sink।
Rate limit-এর জায়গা:
- Network bandwidth: ২৫০ MB/s × ৩ replication = ৭৫০ MB/s। Single 10 GbE NIC ১.২৫ GB/s — fine, কিন্তু margin কম।
- Disk IOPS: NVMe প্রয়োজন। SATA SSD ১৫০ MB/s sustained — bottleneck।
- Single partition throughput: ১০-৫০ MB/s realistic। তাই partition ২৪ চাই।
- Producer batching:
linger.ms=10, batch.size=64KBনা হলে CPU bottleneck। - Consumer lag: downstream slow হলে — disk storage তে event জমে। ২৪ ঘণ্টা lag = ৪ TB extra disk।
Operational concerns:
- Multi-AZ — Bangladesh-এ AWS Mumbai-তে ৩ AZ, এক datacenter fail-এও safe।
- BFIU compliance — sensitive data encrypt at rest (KMS), TLS in transit।
- Data sovereignty — কিছু সংস্থা চায় Bangladesh-এ on-prem Kafka। OracleClouds Bangladesh region option।
- Monitoring: Prometheus + JMX, alert on broker lag, ISR shrink, disk >৭৫%।
মূল উপলব্ধি: Kafka sizing math-এর সাথে judgment। প্রতিটি বড় deployment-এ — load test, headroom (২× peak), disaster recovery practice — design-এ অপরিহার্য।
প্র ০৩ ZooKeeper থেকে KRaft-এ migration একটি big change। কেন Kafka community এই পরিবর্তন করল? Production-এ migration risks কী?
এটি Kafka-র সবচেয়ে বড় architectural shift গত এক দশকে। কারণ বুঝতে — ZooKeeper-এর সমস্যা বুঝতে হবে।
ZooKeeper-এর সমস্যাবলী:
(১) Operational double burden:
- Kafka cluster + আলাদা ZooKeeper cluster — দু'টি system tune করতে হয়।
- Different deployment style — ZK odd number (3, 5), Kafka any number।
- Different security model — SASL, ACLs দু'টি জায়গায়।
- Different monitoring, alerting, debugging tools।
(২) Scalability ceiling:
- ZooKeeper সব partition metadata in-memory রাখে। ২ লাখ partition এ ZK saturate।
- Controller failover slow — ZK থেকে metadata reload-এ মিনিট লাগে large cluster-এ।
- Cell failure-এর সময় rolling restart দীর্ঘ।
(৩) Performance issues:
- ZK ZAB protocol — synchronous, fewer optimization opportunities।
- Watcher fanout overhead — অনেক broker একটি change-এ notify।
- Unnecessary indirection — Kafka-র data Kafka জানে, কিন্তু control plane বাইরে।
(৪) Conceptual mismatch:
- ZK general-purpose coordination service — Kafka-র need-এ tightly fit না।
- Kafka নিজেই একটি replicated log — log-এর জন্য আরেকটি consensus system absurd।
KRaft-এর সুবিধা (KIP-500 vision):
- Million-partition clusters: Metadata একটি Kafka topic হিসেবে stored। In-memory চাপ নেই।
- Faster controller failover: Sub-second (ZK 30s+)।
- Simpler ops: একটি system, একটি deployment।
- Smaller footprint: Embedded — dedicated ZK servers অপ্রয়োজন।
- Future-proof: Kafka 4.0 (২০২৪+) ZooKeeper-free হবে।
Migration risks:
- Maturity: KRaft GA (২০২২) — ZooKeeper ১০+ বছর battle-tested।
- Edge cases: Network partition handling, controller election bugs — early KRaft versions-এ পাওয়া গেছে।
- Tooling lag: kafka-monitor, cruise-control — কিছু tools KRaft-এ slow adopted।
- Migration path: ZK → KRaft live migration KIP-866 সম্প্রতি। Production-এ careful step-by-step।
- Backup/restore: ZK snapshot-based recovery practice ছিল — KRaft-এ নতুন pattern শিখতে।
- Confluent Platform vs OSS: Confluent-এর own additions থাকতে পারে।
Production strategy:
- নতুন cluster — KRaft mode-এ শুরু (now default)।
- পুরনো ZK-based cluster — Kafka 3.6+ stable হলে migrate (Kafka 4.0-এ ZK gone)।
- Test environment-এ আগে — production-এ careful rolling migration।
- Monitoring-এ controller metrics, KRaft replication lag track।
মূল উপলব্ধি: KRaft migration একটি tech debt clearance। Operationally সহজ, scalability বেশি, conceptually pure। ২০২৫-২৬-এর মধ্যে সব Kafka cluster KRaft হবে — ZooKeeper history-তে।
প্র ০৪ "Kafka exactly-once delivery" — এই দাবি কতটা সত্যি? Producer, broker, consumer — তিন স্তরে কোথায় duplicate হতে পারে?
Distributed systems-এর সবচেয়ে nuanced topic। "exactly-once" শব্দটি technically misleading — proper term effectively-once।
Three delivery semantics:
- At-most-once: message একবার বা শূন্যবার। Loss সম্ভব। Performance highest।
- At-least-once: message একবার বা একাধিকবার। Duplicate সম্ভব। Default Kafka behavior।
- Exactly-once: ঠিক একবার। Cross-system-এ অসাধ্য (FLP impossibility)।
Producer স্তরে duplicate:
- Producer broker-এ message পাঠাল, broker write করল কিন্তু network drop-এ ack পৌঁছাল না।
- Producer retry — same message আবার লেখা।
- সমাধান — Idempotent producer:
enable.idempotence=true। Producer-এর প্রতিটি batch-এ producer ID + sequence number। Broker duplicate detect করে drop করে। Single partition-এ exactly-once।
Multi-partition transaction:
- একটি producer একসাথে multiple topic/partition-এ লিখছে — atomic হতে হবে।
- সমাধান:
transactional.id+initTransactions()+commitTransaction()। - Two-phase commit — সব partition success বা সব rollback।
Broker স্তরে duplicate:
- Replication-এ followers leader থেকে message পেয়েছে, কিন্তু leader fail-এর সময় unclean leader election → কিছু message lost বা duplicate।
- সমাধান:
unclean.leader.election.enable=false। ISR-এর বাইরে leader election নয়।
Consumer স্তরে duplicate:
- Consumer message পড়ল, process করল, কিন্তু offset commit করার আগে crash।
- Restart-এ আবার একই offset থেকে পড়ে — duplicate process।
- সমাধান ১ — Idempotent consumer: Output sink-এ idempotent operation। যেমন database-এ
INSERT ... ON CONFLICT (event_id) DO NOTHING। - সমাধান ২ — Transactional consumer + producer: পড়া ও লেখা একই transaction-এ। Kafka Streams নিজে এটি handle করে।
isolation.level=read_committed।
End-to-end exactly-once (Kafka → Kafka):
- Kafka Streams বা Flink-এ
processing.guarantee=exactly_once_v2। - Read offsets + state + write — সব single transaction।
- Performance penalty: ১০-২০% throughput কম, latency বেশি।
External system-এ output (Kafka → MySQL):
- "Exactly-once" cross-system সবসময় effectively-once with idempotency।
- Sink-এ unique key constraint, upsert, deduplication।
- Kafka Connect-এ specific connector (e.g., JDBC sink with PK) handle করে।
Performance trade-off:
- At-most-once: fastest, ~১M+ msg/sec/broker।
- At-least-once: ~৮০০K msg/sec।
- Exactly-once: ~৬০০K msg/sec।
Bangladesh fintech reality:
- bKash transaction-এ at-least-once + idempotent consumer (database PK) — সব duplicate filter।
- Pure exactly-once production-এ rare — overhead too high for most use cases।
- "Once-and-only-once" myth — সবসময় downstream idempotency দরকার।
মূল উপলব্ধি: "Exactly-once" marketing term। সত্যিকার engineering-এ at-least-once + idempotent processing = effectively-once। Kafka-র transactional API powerful, কিন্তু overhead-এর সাথে — শুধু critical path-এ ব্যবহার করুন।
অনুশীলন
-
Partition হিসাব: একটি topic-এ ৬টি partition, ৩টি consumer একটি consumer group-এ। প্রতি consumer কতগুলো partition handle করবে? ১০ম consumer যোগ করলে কী হবে?
- ৬ partition / ৩ consumer = প্রতিজন ২ partition।
- ১০ম consumer যোগ করলে — partition শুধু ৬টি, ৬ consumer max active। বাকি ৪ idle।
- Rule: partition consumer parallelism-এর upper bound। Scale-এ আরও partition দরকার।
- Partition rebalance-এ ১-২ সেকেন্ড "stop the world" — production-এ optimize
cooperative-stickyassignor দিয়ে।
-
Python producer/consumer:
kafka-pythonদিয়ে একটি producer ও consumer লিখুন।from kafka import KafkaProducer, KafkaConsumer import json # Producer producer = KafkaProducer( bootstrap_servers=['localhost:9092'], value_serializer=lambda v: json.dumps(v).encode('utf-8'), key_serializer=lambda k: k.encode('utf-8'), acks='all', enable_idempotence=True ) producer.send('bkash-transactions', key='01711000001', value={'amount': 500, 'type': 'send_money'}) producer.flush() # Consumer consumer = KafkaConsumer( 'bkash-transactions', bootstrap_servers=['localhost:9092'], group_id='fraud-detector', auto_offset_reset='earliest', value_deserializer=lambda v: json.loads(v.decode('utf-8')) ) for msg in consumer: print(f"partition={msg.partition} offset={msg.offset}") print(f" key={msg.key} value={msg.value}") -
Design challenge: Daraz-এ একটি event "order placed" — কোন topic-এ যাবে, কী partition key, কোন consumer group দরকার?
- Topic:
orders(বা domain-prefixed:commerce.orders.v1)। - Partition key:
order_id(uniform distribution) বাcustomer_id(per-user ordering)। - Partitions: peak TPS অনুযায়ী — Daraz-এ ১২-২৪ partition।
- Consumer groups:
inventory-decrementer— stock কমাতে।fraud-checker— high-value order verify।email-sender— confirmation email।analytics-pipeline— real-time dashboard।warehouse-router— কোন warehouse থেকে dispatch।s3-archiver— long-term storage।
- Schema: Avro বা JSON Schema with Schema Registry — backward compatible evolution।
- Topic:
আরও পড়ুন · ABCL TECH-এ আপনার পরবর্তী পদক্ষেপ
- পাঠ ১৮ · Producer, Consumer, Topic পরবর্তী পাঠ Kafka API-র গভীরে — producer config, consumer group, offset management।
- পাঠ ১৬ · Streaming vs Batch আগের পাঠ Streaming-এর জগৎ-এ ঢোকার আগের সিদ্ধান্ত — কখন কোনটি দরকার।
- পাঠ ২১ · Change Data Capture (CDC) এই পাঠের সাথে সম্পর্কিত Database থেকে Kafka-তে real-time sync — Debezium ও Kafka Connect।
- সব AI Courses দেখুন ABCL TECH Python, ML, DL, NLP, CV, GenAI, RL, MLOps — সব AI কোর্স একসাথে।