Skip to main content

Posts

Showing posts with the label Event Streaming

Kafka on Object Storage Was Inevitable. The Next Step Is Open.

Kafka on object storage is not a trend. It is a correction. WarpStream proved that the Kafka protocol can run without stateful brokers by pushing durability into object storage. The next logical evolution is taking that architecture out of vendor-controlled control planes and making it open and self-hosted. KafScale is built for teams that want Kafka client compatibility, object storage durability, and Kubernetes-native operations without depending on a managed metadata service. The problem was never Kafka clients The Kafka protocol is one of the most successful infrastructure interfaces ever shipped. It is stable, widely implemented, and deeply integrated into tooling and teams. The part that aged poorly is not the protocol. It is the original broker-centric storage model. Stateful brokers made sense in a disk-centric era where durability lived on the same machines that ran compute. That coupling forces partition rebalancing, replica movement, disk hot spots, slow recovery, an...

Why Your Kafka Streams App Uses 5x More Memory Than You Configured

Kafka Streams uses RocksDB for stateful operations, creating hundreds of database instances per node. Each instance allocates native memory outside JVM heap tracking. A 4GB heap application can consume 20GB+ actual RAM. This article covers production evidence of memory issues, 24x write amplification killing SSDs, compaction storms causing latency spikes, and tuning parameters. Alternatives to RocksDB include in-memory stores, custom StateStore implementations, and remote state backends. Every Kafka Streams application with stateful operations uses RocksDB by default. Most teams discover this fact only after their production system starts exhibiting mysterious latency spikes, memory pressure, or disk I/O that saturates their SSDs. I have debugged RocksDB-related Kafka Streams issues across multiple production deployments. The pattern is consistent: teams scale their Kafka infrastructure, tune their consumer configurations, optimize their topology, and still hit a wall. The wall is...

Stream IoT data to S3 - the simple way

This article introduces infinimesh as a Kubernetes-native IoT platform built to integrate massive device fleets without cloud lock-in, and highlights its expanding plugin ecosystem for Elastic, Redis TimeSeries, SAP HANA, Snowflake and cloud-native object storage. Using lightweight Go-based Docker containers, plugins can run securely in user environments—even on AWS free-tier instances—while streaming device data directly into existing infrastructures like S3 or MinIO. The step-by-step example shows how quickly CloudConnect can move IoT data into object storage using docker-compose, with the plugin internally batching device payloads through Redis before exporting them as CSV. In benchmarking, the architecture handled millions of devices sending frequent JSON updates, demonstrating how infinimesh plugins simplify large-scale IoT data integration with minimal resource overhead. First, a short introduction to infinimesh , an Internet of Things (IoT) platform which runs completely in Kub...

Handling Corrupted Kafka Messages and Offset Recovery in Distributed Systems

This article explains how corrupted Kafka messages occurred in early Kafka versions, how offsets were stored in Zookeeper and how to manually recover a stuck consumer. It documents the race condition described in KAFKA 2477, shows how to inspect offsets using Kafka tools or Zookeeper and describes code based and operational strategies for skipping bad messages in older distributed log systems. Handling Corrupted Kafka Messages and Offset Recovery in Distributed Systems In older Kafka deployments, especially versions before 0.9, it was possible for a message in a topic to become unreadable due to corruption. This happened most often when third party frameworks interacted with Kafka internals or when the consumer logic encountered a rare race condition. One such condition was documented in KAFKA 2477, where a lock on Log.read was missing at the consumer level while Log.write remained protected. Under specific timing, this resulted in a corrupted message being written to ...