Direct Zero-ETL Ingestion: Streaming Apache Kafka into ClickHouse via Kafka Engine & Materialized Views

Eliminate brittle intermediate consumer microservices. Learn how to stream 200,000 Kafka events/sec directly into ClickHouse using the native Kafka Engine.

The Overhead of Traditional ETL Consumer Microservices

Streaming high-throughput event queues into analytical data warehouses traditionally requires dedicated ETL ingestion clusters: Logstash pods, Spark Streaming jobs, or custom Python Celery consumers. While functional, these intermediate consumer layers introduce severe operational liabilities: consumer lag debugging, serialized JSON re-encoding overhead, worker node memory leaks, and complex distributed transaction handling.

For organizations operating at hundreds of thousands of events per second, ClickHouse offers a vastly superior architecture: the native Kafka Engine Table combined with background Materialized Views. In this setup, ClickHouse itself acts as a multi-threaded Kafka consumer group, batching records in memory and atomically committing them into columnar MergeTree storage with zero intermediary infrastructure.

1. Architecture: The 3-Tier Zero-ETL Pipeline

The native ClickHouse ingestion pipeline consists of three coordinated primitives:

  1. The Kafka Engine Table: A virtual buffer that subscribes to Kafka topics, reads serialized byte payloads, and commits consumer group offsets.
  2. The Destination Columnar Table: A high-performance ReplacingMergeTree or AggregatingMergeTree table that stores compacted columnar data on disk.
  3. The Materialized View: The reactive trigger pipeline that automatically parses JSON fields from the Kafka table and flushes columnar blocks into the destination table.

2. Production ClickHouse Table DDL Configuration

Here is the complete production configuration for ingesting telemetry events:

-- 1. Destination storage table (compressed, partitioned columnar storage)
CREATE TABLE telemetry.events_stream (
    event_id UUID,
    tenant_id LowCardinality(String),
    event_name LowCardinality(String),
    payload String,
    latency_ms Float32,
    timestamp DateTime64(3, 'UTC')
) ENGINE = ReplacingMergeTree()
PARTITION BY toYYYYMM(timestamp)
ORDER BY (tenant_id, event_name, timestamp, event_id);

-- 2. Virtual Kafka consumer table
CREATE TABLE telemetry.events_kafka_queue (
    event_id UUID,
    tenant_id String,
    event_name String,
    payload String,
    latency_ms Float32,
    timestamp DateTime64(3, 'UTC')
) ENGINE = Kafka
SETTINGS 
    kafka_broker_list = 'kafka-1.internal:9092,kafka-2.internal:9092',
    kafka_topic_list = 'production.telemetry.events',
    kafka_group_name = 'clickhouse_telemetry_consumers',
    kafka_format = 'JSONEachRow',
    kafka_max_block_size = 65536,
    kafka_num_consumers = 4;

-- 3. Continuous ingestion Materialized View
CREATE MATERIALIZED VIEW telemetry.events_mv TO telemetry.events_stream AS
SELECT 
    event_id,
    tenant_id,
    event_name,
    payload,
    latency_ms,
    timestamp
FROM telemetry.events_kafka_queue;

3. Tuning Batch Sizes to Prevent Small-Part Fragmentation

ClickHouse achieves legendary query performance by compacting massive blocks of data onto disk. If the Kafka engine flushes every single message independently, ClickHouse will quickly crash with the dreaded Too many parts error. To enforce massive block ingestion, configure these parameters in ClickHouse's config.xml:

  • kafka_max_block_size = 65536: Buffers up to 64k records in memory before writing a single physical part to disk.
  • stream_flush_interval_ms = 1000: Flushes data at least once per second to maintain near-real-time visibility.

Pairing this zero-ETL ingestion pipeline with our guide on ClickHouse AggregatingMergeTree Sizing provides cost-effective analytics over billions of events. Explore our Data Pipelines & ETL Consulting.

Interactive PostgreSQL Memory & Tuning Calculator

// Real-Time Production Memory Allocator
PostgreSQL 14 / 15 / 16 / 17

Adjust your server resources below to calculate optimized postgresql.conf memory thresholds, autovacuum scale factors, and cost weights.

16 GB
2 GB 64 GB 256 GB
8 Cores
2 16 64
100
20 200 1,000
generated-postgresql.conf
# Memory Allocations
shared_buffers = 4GB
work_mem = 40MB
maintenance_work_mem = 1GB
effective_cache_size = 12GB

# Concurrency & Background Workers
max_connections = 100
max_worker_processes = 8
max_parallel_workers_per_gather = 4
max_parallel_workers = 8

# Autovacuum Tuning (Prevent Bloat)
autovacuum_max_workers = 4
autovacuum_vacuum_scale_factor = 0.05
autovacuum_analyze_scale_factor = 0.02
autovacuum_vacuum_cost_limit = 1000

# Planner Cost Constants (NVMe SSD)
random_page_cost = 1.1
effective_io_concurrency = 200
Need hands-on database profiling? We analyze query execution plans, resolve lock trees, and eliminate replication lag.
Book Database Audit (30m)
// Production Systems Architecture • Database Diagnostic Audit

Diagnosing Production PostgreSQL Bloat, Lock Contention, or Replication Lag?

Theoretical tuning only goes so far. We provide hands-on architectural reviews of query execution plans, autovacuum parameters, connection pools, and read-replica lag for high-concurrency systems.

All Insights
Chat on WhatsApp