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:
- The Kafka Engine Table: A virtual buffer that subscribes to Kafka topics, reads serialized byte payloads, and commits consumer group offsets.
- The Destination Columnar Table: A high-performance
ReplacingMergeTreeorAggregatingMergeTreetable that stores compacted columnar data on disk. - 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.