Skip to main content

Posts

Showing posts with the label Stream Processing

Production CDC Architecture: Debezium Scaling Lessons

Production CDC architecture breaks under load long before most teams expect it. With Debezium, Kafka Connect, and Postgres, the failure patterns are consistent: WAL pressure builds up, connector lag drifts unnoticed, and snapshot phases exhaust memory under bursty traffic. This is based on running these pipelines across high throughput systems, including workloads above 10k TPS. The difference between a system that works and one that holds under pressure comes down to observability, WAL discipline, and how connector scaling is handled. Production Debezium CDC Architecture Operational reality vs. tutorial defaults under real load (10k+ TPS) The Default "Tutorial" Setup Assumes low throughput and stable networks. Fails under pressure. Source: Postgres Single WAL Slot Shared slot coupling multiple connectors Default WAL retention settin...

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...

Building Reliable Flink-to-Iceberg Pipelines for Unity Catalog and Snowflake

Apache Flink ®, Apache Iceberg ® and governed catalogs such as Databricks Unity Catalog or Snowflake are often pitched as a simple path from Apache Kafka ® JSON to managed tables. In reality Flink is a stream processor, Iceberg is an open table format and the catalog handles governance. None of them infers schemas or models messy payloads for you. You still design schemas, mappings and operations under real Java, DevOps and cost constraints. Many architectural diagrams show a clean pipeline: Kafka into Flink, Flink into Iceberg, Iceberg governed by Unity Catalog or queried from Snowflake. In practice this stack has real friction. Flink is not a neutral glue layer. It is a JVM-centric stream processor with non-trivial operational cost. Iceberg is not a storage engine but a table format that imposes structure. Unity Catalog and Snowflake add their own expectations around governance and schema. Apache Flink is a distributed stream processor for stateful event pipelines. Apache Iceberg i...

SolR, NiFi, Twitter and CDH 5.7

Since the most interesting Apache NiFi parts are coming from ASF [1] or Hortonworks [2], I thought to use CDH 5.7 and do the same, just to be curious. Here's now my 30 minutes playground, currently running in Googles Compute. On one of my playground nodes I installed Apache NiFi per mkdir /software && cd /software &&  wget http://mirror.23media.de/apache/nifi/0.6.1/nifi-0.6.1-bin.tar.gz   && tar xvfz nifi-0.6.1-bin.tar.gz Then I've set only nifi.sensitive.props.key property in conf/nifi.properties to an easy to remember secret. The next bash /software/nifi-0.6.1/bin/nifi.sh install installs Apache NiFi as an service. After log in into Apache NiFi's WebUI, download and add the template [3] to Apache NiFi, move the template icon to the drawer, open it and edit the twitter credentials to fit your developer account. To use an  schema-less SolR index (or Cloudera Search in CDH) I copied some example files over into a local directory: cp -r ...

Practical Memory Sizing for Apache Flume Sources, Sinks and File Channels

Apache Flume still appears in many legacy data estates, and most operational issues come from undersized heap or direct memory. This updated guide explains how to estimate memory requirements for Flume sources, sinks and file channels, how batch sizing impacts heap usage, and how replay behavior can drastically increase memory demand. The goal is to give operators a reliable sizing baseline instead of trial-and-error tuning. Memory Requirements for Flume Sources and Sinks The dominant memory cost for each event comes from its body plus a small overhead for headers (typically around 100 bytes, depending on the transaction agent). To estimate memory for a batch: Take the average or p90 event size. Add a buffer for headers and variability. Multiply by the maximum batch size. This result approximates the memory required to hold a batch in a Source or Sink. A Sink needs memory for one batch at a time. A Source needs memory for one batch multiplied by the number of ...

Flume 1.2.0 released

The Apache Flume Team released yesterday the next large release with number 1.2.0. Here a overview about the fixes and additions (thanks Mike, I copy your overview): Apache Flume 1.2.0 is the third release under the auspices of Apache of the so-called "NG" codeline, and our first release as a top-level Apache project! Flume 1.2.0 has been put through many stress and regression tests, is stable, production-ready software, and is backwards-compatible with Flume 1.1.0. Four months of very active development went into this release: a whopping 192 patches were committed since 1.1.0, representing many features, enhancements, and bug fixes. While the full change log can be found in the link below, here are a few new feature highlights: * New durable file channel  * New client API  * New HBase sinks (two different implementations)  * New Interceptor interface (a plugin processing API)  * New JMX-based monitoring support With this release - the first after e...

Using the Apache Flume HBase Sink: How the Integration Works and How to Configure It

The first Apache Flume HBase sink introduced a simple way to stream events directly into HBase tables. This modernized walkthrough explains how the sink works, what its limitations are, how Flume resolves HBase configuration files, and how to set up a minimal but functional Flume-to-HBase pipeline. Although this feature originated in early Flume versions, many legacy Hadoop deployments still rely on it today. Overview The HBase sink was added to the Flume trunk and provided direct write support from Flume channels into HBase tables. It relies on synchronous HBase client operations and requires that HBase table metadata already exists. The sink handles flushes, transactions and rollbacks, allowing Flume to treat HBase as a durable storage target. Building Flume from Trunk In early versions the HBase sink was only available in the trunk source. The following sequence checks out Flume and builds it using Maven: git clone git://git.apache.org/flume.git cd flume git checkout tr...

Getting Started with Apache Flume NG: Flows, Agents and Syslog-to-HDFS Examples

Apache Flume NG replaced the original master/collector architecture with lightweight agents that can be wired together to form flexible data flows. This guide explains what changed with Flume NG, how the agent–channel–sink model works, and walks through simple configurations for syslog ingestion to a console logger and to HDFS. It’s aimed at engineers who still operate Flume in legacy estates or need to understand it for migrations. From Flume to Flume NG Apache Flume is a distributed log and event collection service. With Flume NG, the project moved away from the original master/client and node/collector design and adopted a simpler, more robust architecture based on standalone agents. Key changes introduced by Flume NG: No external coordination service required for basic operation. No master/client or node/collector roles—only agents . Agents can be chained together to build arbitrary flows and fan-in/fan-out patterns. Lightweight runtime; small heap sizes are s...