Distributed Systems

Distributed Systems & Event-Driven Architecture

Distributed systems fail in combinations: a slow dependency, duplicated event, partial write or retry storm can cross boundaries quickly. I design the failure path before trusting the happy path.

My approach

The moment a business transaction crosses process or network boundaries, failure becomes partial. One service may succeed while another times out. A message may be delivered more than once. A consumer may restart after writing to a database but before acknowledging an event.

I therefore avoid designs that depend on exactly-once behavior across independent systems unless the underlying platform genuinely provides it for the complete transaction. Most real systems become safer when duplication, retries and delayed processing are expected and handled explicitly.

Event-driven architecture is powerful when it reduces coupling and supports asynchronous workflows, but Kafka is not a substitute for domain boundaries, data ownership or operational discipline. Topics, keys, schemas, consumer behavior and recovery all need architecture decisions.

Capability

What I focus on

I prefer to describe expertise through architecture decisions and production responsibilities rather than a list of tools.

01

Kafka and event design

Define event ownership, keys, ordering, schema evolution, retention and replay so the event stream remains useful after the first implementation.

02

Idempotency and duplicate handling

Assume at-least-once delivery where appropriate and make repeated commands/events safe through keys, state checks or deduplication records.

03

Backpressure and load shedding

Protect slower downstream systems using bounded concurrency, consumer pause/resume, queues, rate limits and deliberate degradation.

04

Transactional outbox

Keep database state and event publication consistent without relying on fragile dual writes across unrelated transactional systems.

05

Saga and workflow coordination

Use orchestration or choreography with explicit compensations when business workflows span independently committed services.

06

Resilience patterns

Apply timeouts, retry budgets, exponential backoff, jitter, circuit breakers and bulkheads based on dependency behavior.

Architecture

A durable event-driven transaction flow

The critical idea is to commit local state once, publish reliably, and make downstream processing repeatable. Recovery should be possible without manual reconstruction of hidden state.

API / Command
→
Local DB Transaction
→
Outbox
→
Kafka
→
Idempotent Consumer
→
Downstream State

Local atomicity

Keep the business state change and outbox record in the same database transaction.

At-least-once safe

Consumers should tolerate duplicate delivery and restart without corrupting business state.

Replayable

Retain enough event history, schema discipline and consumer design to support recovery and new consumers.

Production practice

How I approach production design

01

Design events as contracts, not log messages

An event should communicate a meaningful fact with stable identity and ownership. I avoid publishing internal database rows as public contracts because that couples consumers to implementation details.

Keys matter because they define partitioning and ordering. Schema evolution matters because producers and consumers rarely deploy at exactly the same time.

  • Event identity
  • Aggregate/business key
  • Schema versioning
  • Producer ownership
  • Retention
  • PII classification
02

Make duplicate processing safe

A consumer can process successfully and fail before acknowledging the message. That same message can then arrive again. If the handler charges a card, creates inventory or sends a notification twice, the system has a business problem—not merely a messaging problem.

Idempotency can be implemented through natural business keys, state transitions, deduplication tables or conditional writes depending on the operation.

03

Treat retries as additional load

Retries are useful for transient failures, but every retry consumes resources while the dependency is already unhealthy. I define retryable error classes, attempt limits, backoff and jitter, and avoid stacking retries at gateway, service and client layers.

When work can be delayed, retry topics or queues may provide better protection than holding synchronous request threads open.

04

Use sagas for business recovery, not technical rollback

A distributed business workflow cannot always be rolled back like one SQL transaction. Compensation is a new business action: refund payment, release reservation, cancel shipment or mark an order for manual handling.

I document compensating actions and states explicitly so operations teams can understand partially completed workflows.

Decision framework

Questions I want answered before approving the design

API or event?

Use synchronous APIs when the caller requires an immediate outcome. Use events when consumers can react independently, processing can be asynchronous or multiple downstream capabilities need the same business fact.

What should the Kafka partition key be?

Choose the key based on the ordering boundary the business actually needs. A poor key can create hot partitions or ordering assumptions that cannot be maintained.

Outbox or direct publish?

Use an outbox when database state and event publication must not diverge. Direct publishing may be acceptable when there is no coupled local transaction or when the platform provides an equivalent reliable mechanism.

Orchestration or choreography for a saga?

Orchestration makes workflow state and recovery explicit for complex processes. Choreography can keep simple reactions decoupled but becomes hard to reason about when many services implicitly form one business workflow.

When should a message go to a DLQ?

After bounded retries when the failure is unlikely to succeed automatically. A DLQ also needs ownership, diagnostics, replay tooling and business rules; otherwise it becomes a permanent error archive.

Reliability

Failure, scale and operational reality

Risk

Consumer faster than downstream database

Architecture response

Bound concurrency, monitor lag, pause/resume consumption where appropriate, batch carefully and scale the constrained resource rather than simply adding more consumers.

Risk

Duplicate event delivery

Architecture response

Use stable event/operation identifiers and idempotent handlers. Test replay and restart behavior as normal scenarios.

Risk

Poison message

Architecture response

Classify deterministic vs transient failure, use bounded retries, send unrecoverable items to a managed DLQ path and preserve enough context for replay.

Risk

Schema change breaks consumers

Architecture response

Use backward/forward-compatible evolution rules, schema validation and staged migrations instead of synchronized producer/consumer releases.

Risk

Retry storm

Architecture response

Centralize retry policy, use exponential backoff and jitter, cap attempts, apply circuit breakers and reduce traffic while dependencies recover.

Security

Security by architecture

  • Authenticate and authorize producers/consumers; topic access should follow least privilege.
  • Classify event payloads for PII and secrets before choosing retention and replication policies.
  • Encrypt traffic and protect credentials used by brokers, schema registries and clients.
  • Avoid putting data into events merely because it is convenient; events are often retained and copied widely.
  • Audit sensitive workflow transitions and administrative replay operations.

Observability

Operate what we design

  • Monitor consumer lag by consumer group and partition, not only aggregate throughput.
  • Track publish/consume latency, retry count, DLQ volume and handler duration.
  • Correlate events using business identifiers and trace context where useful.
  • Alert on stalled partitions, hot keys and growing outbox backlog.
  • Measure downstream saturation so backpressure decisions are based on the actual bottleneck.

Continue reading

Related architecture guides

FAQ

Frequently asked questions

Does Kafka guarantee exactly-once processing?

Kafka provides exactly-once capabilities in specific Kafka transaction/streaming scenarios, but end-to-end business processing that includes external databases or APIs still needs careful transactional and idempotency design.

Why are Kafka consumers commonly designed to be idempotent?

Because at-least-once delivery and consumer restart behavior can result in the same message being processed more than once. Idempotency makes repeated delivery safe.

What problem does the transactional outbox solve?

It avoids the dual-write problem where a database commit succeeds but event publication fails, or vice versa. Business state and an outbox record are committed atomically, then the event is published reliably from the outbox.

What is backpressure?

Backpressure is controlling incoming work when a downstream dependency cannot keep up. In event systems this can involve bounded concurrency, consumer pause/resume, queue depth, batching or load shedding.

When should I use a saga?

Use a saga when one business workflow spans multiple independently committed services and you need explicit progression, compensation and recovery instead of one distributed ACID transaction.

About the author

Romharshan Singh

Senior Solution Architect • AI & Cloud Mentor

I write about architecture from a production perspective: how systems fail, how design decisions affect cost and operability, and how teams can turn technology choices into maintainable enterprise platforms.