At small volumes, streaming data is straightforward: spin up a consumer, parse JSON, and dump into a relational database. But when you cross 100,000 sustained messages per second—peaking significantly higher—every assumption in standard data architecture breaks down.

In my recent engineering work, we designed and operated a distributed real-time platform capable of continuously ingesting and transforming over 9 billion Kafka messages per day, representing approximately 30 TB of operational payload daily.

The business requirement was uncompromising: customer-facing analytics dashboards could no longer wait for hourly or daily batch ETLs. Users needed live operational visibility. If an order changed state, an inventory count updated, or an anomaly occurred, the downstream analytics layer needed to reflect that change within seconds—without locking tables, ballooning cloud costs, or causing query timeouts for concurrent users.

Here is how we designed the architecture, the production trade-offs we navigated, and how we solved the hardest problem in streaming: incremental real-time merging at scale.

The Core Architectural Dilemma: Ingestion vs. Merging

Most teams fail at this scale because they treat streaming ingestion and downstream consumption as the same problem.

Writing 30 TB of raw append-only data into an object store (like S3 or ADLS Gen2) is relatively simple with distributed stream processing. The catastrophic bottleneck occurs when you attempt to merge continuous updates (UPSERTs) into dimensional tables while thousands of analytical queries are hitting the same dataset.

[Source Systems & APIs] │ ▼ (9B+ events / 30 TB daily) [Apache Kafka Cluster: Partitioned by Business Entity Key] │ ├──► [Stream Ingestion & Serialization: Avro + Schema Registry] │ ▼ [Distributed Stream Processing Engine: Stateful Windowing & Deduplication] │ ├──► [Raw Bronze Append Store: Fast, immutable micro-batches] │ ▼ [Streaming Incremental Merge Layer: Compaction + ACID Delta Merges] │ ▼ [Real-Time Analytics Store / Lakehouse (Silver & Gold Marts)] │ ▼ [Customer-Facing Dashboards & Downstream AI Applications]

Key Engineering Decisions That Made It Work

1. Partitioning by Stable Business Entity Key

When processing billions of stateful updates, out-of-order events are inevitable. If customer order #1042 has a "Created" event, followed milliseconds later by an "Updated" event, routing those two events to different Kafka partitions guarantees race conditions and corrupted final states.

We enforced strict partitioning keyed by the root business entity identifier. This ensured that all lifecycle events for a specific transaction arrived in sequential order within the same Kafka partition, allowing consumer workers to maintain deterministic local state without cross-node synchronization locks.

2. Taming the Small-Files Problem in Streaming Merges

A major hazard when streaming into table formats like Delta Lake or Apache Iceberg is the "small file disease." If you commit files every 5 seconds at high throughput, you generate tens of thousands of tiny parquet files an hour, crippling query performance for downstream BI tools.

We solved this with a decoupled two-tier writing pattern:

  • Tier 1 (Fast Bronze Append): Streaming workers write micro-batches directly to an append-only append log. No upserts or heavy file scans happen in this critical ingest loop.
  • Tier 2 (Continuous Incremental Compaction & Merge): A secondary worker pool reads incremental change-data checkpoints, buffers updates in memory using time-based windowing, and performs coordinated partition-pruned merges into the queryable layer while background compaction combines smaller files into optimal 128 MB–256 MB chunks.

Key Takeaway: Never perform full-table ACID merges directly inside the primary Kafka consumer loop. Decouple ingestion from transactional compaction to maintain sub-second consumer lag under burst traffic.

3. Managing Backpressure and Memory Spikes

Network partitions or downstream storage latency can instantly back up consumer lag. At 100,000+ messages per second, a 2-minute delay means 12+ million events queue in memory.

We tuned Kafka consumer fetch limits (`max.poll.records`) and utilized reactive backpressure mechanisms within the stream workers. If downstream write operations encountered storage throttling, the consumer dynamically throttled pull rates rather than crashing workers with out-of-memory (OOM) errors.

4. Schema Evolution Without Pipeline Downtime

In an enterprise environment producing 30 TB a day across multiple upstream engineering teams, schemas change frequently. Hard-coded JSON schemas break pipelines daily.

We enforced Confluent Schema Registry with strict backward and forward compatibility rules using Apache Avro. When upstream applications added non-breaking fields, downstream streaming processors automatically adapted without restarting clusters or interrupting analytics consumers.

The Outcome

By focusing on architecture rather than simply throwing more compute nodes at the problem, the platform achieved:

  • Continuous Freshness: Customer analytics refreshed with sub-minute latency instead of 24-hour batch delay.
  • Predictable Cloud Costs: Efficient compaction and partition pruning eliminated runaway compute costs associated with continuous naive full-table scans.
  • Operational Resilience: Zero unrecoverable pipeline outages during traffic spikes exceeding 150% of typical daily volume.

Facing a high-throughput data or streaming challenge?

We design and modernize production data platforms on Azure, GCP, and Kafka that scale reliably.

Discuss Your Architecture