Skip to content
~/mahadi hassan
← All posts

Running an In-House Kafka Pipeline for a SaaS Ecosystem

Why an interconnected SaaS suite needed its own event backbone, and how I ran a self-hosted Kafka pipeline — topic design, consumer groups, idempotency, and the failure modes you only learn in production.

3 min read
~/mahadi/blog

Running an In-House Kafka Pipeline for a SaaS Ecosystem

Apache KafkaMicroservices
On this page

At Amharc Tech the product was not one app but an ecosystem — system, SaaS and service admin portals, a Web POS, and customer and technician mobile apps, all on a NestJS microservices backend. Once several services need to react to the same business event, point-to-point HTTP calls turn into a spider web. We pulled an in-house Apache Kafka pipeline into the middle of it as the event backbone. Here is what that actually took.

Why an event log, not more HTTP

A booking created in one service ripples outward: notify the customer, allocate a technician, update the calendar, bump analytics. Doing that with synchronous calls couples the booking service to four others and makes it as slow and as fragile as the slowest of them. With Kafka, the booking service appends one event to a topic and forgets about it. Each interested service reads at its own pace, and adding a fifth consumer later changes nothing upstream.

Topic design is the real design

The mistake is to model topics after services. Model them after events that happened — past tense, immutable facts:

// one topic per domain event, keyed by the aggregate id
await producer.send({
  topic: 'booking.created',
  messages: [{ key: booking.id, value: JSON.stringify(event) }],
});

Keying by the aggregate id (the booking id) matters: Kafka guarantees ordering within a partition, and same-key messages always land on the same partition. So every event for one booking is processed in order, even while different bookings fan out across partitions for throughput.

Consumer groups = free horizontal scale

Each service joins with its own consumer group. Kafka hands every group the full stream, and within a group it splits partitions across instances. Run three replicas of the notifications service and Kafka rebalances the partitions across them automatically — no coordination code on my side.

The two rules that keep it sane

  • Consumers must be idempotent. Kafka is at-least-once; a rebalance or retry will redeliver. Every handler dedupes on the event id before acting, so processing the same booking.created twice sends one email, not two.
  • A poison message must not block the partition. A row that always throws will halt everything behind it. Failed messages go to a retry topic with a delay, and after N attempts to a dead-letter topic for inspection — the main partition keeps moving.
async function handle(event: BookingCreated) {
  if (await seen.has(event.id)) return; // idempotent
  try {
    await notify(event);
    await seen.add(event.id);
  } catch (err) {
    await republishWithBackoff(event); // retry topic → DLQ
  }
}

What I would weigh before reaching for it

Kafka is operational weight — brokers, partitions, consumer lag to watch. For a single service it is overkill; a database-backed queue is plenty. It earns its keep the moment multiple independent services react to the same events and you want them decoupled and independently scalable. That was exactly the shape of the ecosystem, and the event log is what kept it from collapsing into a mesh of HTTP calls.

Comments

Loading comments…