পাঠ ১৭ · ২৯-এর মধ্যে · মডিউল ৩

Apache Kafka পরিচিতি

Apache Kafka — the distributed log
৮ মিনিট পড়া মাঝারি · Intermediate Distributed systems

এই পাঠে যা শিখবেন

  • 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-এর ৮০%+ ব্যবহার করে।

Log বনাম Queue — মৌলিক পার্থক্য

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।
ভাবুন একটি বিশাল খবরের কাগজ archive। প্রতিদিন নতুন পাতা যোগ হয় (append-only)। অনেক পাঠক (consumer) একই পাতা একই সময়ে পড়তে পারে। কেউ ১ দিন আগের, কেউ ১ সপ্তাহ আগের — সবার নিজের bookmark (offset)। Kafka ঠিক এমনই — কিন্তু খবর নয়, business events; এবং পাতা ১০টি কাগজে duplicate (replication)।

৩ · 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)।

ISR (In-Sync Replicas): যে followers leader-এর সাথে up-to-date — তারা ISR। Producer যদি 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।
Kafka Cluster — topic বিতরণ topic: bkash-transactions, partitions=3, RF=3 📱 Producers bKash mobile app Agent POS terminal Web dashboard ⚙ Kafka Cluster (3 brokers) Broker 1 P0 leader P1 follower P2 follower KRaft node Broker 2 P0 follower P1 leader P2 follower KRaft node Broker 3 P0 follower P1 follower P2 leader KRaft controller leader = green, follower = yellow replication factor = 3 👥 Consumer groups fraud-detector offset: 1,254,300 analytics-pipeline offset: 980,123 archiver-to-s3 offset: 1,254,290 প্রতি group আলাদা offset Producer publish → Kafka durable log → অনেক consumer pull (নিজ গতিতে) Decoupled · Durable · Replayable
৩-broker Kafka cluster — topic-এর ৩ partition প্রতিটি broker-এ leader, ৩×replication। Producer ও consumer সম্পূর্ণ decoupled।

৭ · 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):

YAML · docker-compose.yml
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 দিয়ে:

Bash · Kafka 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

    
Producer-এ লেখা message instantly consumer-এ আসবে। --from-beginning মানে topic-এর শুরু থেকে — Kafka-র replay capability প্রদর্শিত।
Production-এ যা মনে রাখবেন: single broker = single point of failure। কমপক্ষে ৩-broker cluster, replication factor ৩, 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:

  1. নতুন cluster — KRaft mode-এ শুরু (now default)।
  2. পুরনো ZK-based cluster — Kafka 3.6+ stable হলে migrate (Kafka 4.0-এ ZK gone)।
  3. Test environment-এ আগে — production-এ careful rolling migration।
  4. 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-এ ব্যবহার করুন।

অনুশীলন

  1. 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-sticky assignor দিয়ে।
  2. 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}")
  3. 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।

আরও পড়ুন · ABCL TECH-এ আপনার পরবর্তী পদক্ষেপ

কোড রানার কাজ না করলে? Kafka local-এ install কঠিন হলে Google Colab বা Confluent Cloud-এর free tier ব্যবহার করুন। Docker desktop থাকলে docker-compose সবচেয়ে সহজ।
পূর্ববর্তী পাঠ
পাঠ ১৬ · Streaming vs Batch