Q: Your microservices fleet generates 100 Terabytes of unstructured and JSON logs every single day (~1.2 million events/sec). Commercial SaaS tools (Datadog, Splunk) would cost over $250,000/month. How do you architect an in-house, distributed, highly available log aggregation and indexing platform that provides sub-2-second search queries while keeping infrastructure costs under $15,000/month?
Architecting an ultra-high-throughput, cost-efficient observability logging pipeline capable of ingesting 100 TB/day (1.2 million log events/sec) with sub-second query latency using Vector, Kafka, ClickHouse, and object storage tiering.
Want to master this scenario in a live sandbox? The Linux Foundation's FinOps Certified Practitioner (FOCP) Program covers this exact problem with hands-on terminal drills.
🛠️ Production Runbook & Step-by-Step Resolution
Deploy High-Performance Edge Log Collectors with Vector DaemonSets
Collect, parse, and compress log streams at the node layer without taxing host CPU:
- Vector DaemonSets: Deployed Vector (written in Rust) on every Kubernetes worker node, consuming stdout logs from container runtimes via
/var/log/pods. - Schema Normalization: Vector parses JSON payloads, drops useless health check logs, enriches logs with Kubernetes metadata (pod, namespace, git commit), and compresses payloads using Zstandard (zstd level 3).
Buffer Ingestion Bursts with Apache Kafka / Redpanda Streaming Tier
Decouple high-frequency ingest spikes from downstream indexing pipelines:
- Partitioning Strategy: Provisioned a 12-broker Kafka cluster with 96 partitions distributed across 3 Availability Zones, partitioned by
tenant_idorservice_name. - Retention & Throughput: Configured 6-hour disk retention buffer, absorbing massive downstream ClickHouse maintenance restarts or query spikes without dropping a single log entry.
Store & Index Columnar Logs with ClickHouse Distributed Engine
Ingest structured events into ClickHouse for instantaneous SQL querying:
- Kafka Engine Table: ClickHouse consumer tables consume directly from Kafka partitions in micro-batches (50,000 rows or 2 seconds).
- MergeTree Partitioning: Data is written to
ReplacingMergeTreepartitioned bytoYYYYMMDD(timestamp)and ordered by(tenant_id, service_name, level, timestamp). - Compression: ClickHouse columnar encoding achieves 8:1 to 10:1 data compression ratios on log data.
Implement Tiered Cold Storage on Cloud Object Storage (S3 / GCS)
Optimize multi-petabyte storage expenses by offloading historical data to object storage:
- Hot Tier (NVMe SSD): Retains 3 days of raw logs on local high-speed NVMe disks for instant incident triaging.
- Cold Tier (S3 / GCS): ClickHouse Storage Policy automatically moves parts older than 3 days to AWS S3 / Google Cloud Storage, retaining 90 days of searchable historical records.
- Total Cost: S3 storage + 12 ClickHouse compute nodes cost $11,200/month, slashing monitoring expenses by 95% compared to commercial SaaS.
- Collect and enrich logs at the node layer using high-performance Vector DaemonSets.
- Buffer bursty traffic through a multi-partition Kafka cluster to guarantee zero data loss.
- Ingest micro-batches into ClickHouse columnar storage for sub-second analytical search.
- Tier logs older than 3 days to S3/GCS object storage to maintain 90-day retention under budget.