Introduction
Countly has grown with the needs of modern analytics, expanding far beyond what the original Node.js and MongoDB ingestion architecture could sustainably deliver. As customers began sending billions of data points per month and expecting real-time analytics at greater precision, Countly required a new foundation.
This document explains the next-generation architecture in a narrative, approachable way while preserving the technical depth that engineering teams expect.
It walks the reader through the reasoning behind the redesign, its components, its operational reality, and what migration feels like from a customer’s point of view.
Why Countly Evolved: The Shift to a Streaming-Based Architecture
Countly’s legacy architecture was optimized for simplicity:
SDKs sent data → Node.js ingestion → MongoDB stored events → dashboards queried aggregated collections.
It worked extraordinarily well for years until customers began growing into tens of billions of monthly events.
MongoDB became the bottleneck. Not because MongoDB is slow, but because the use case required a form of elastic throughput and analytical scanning that row-store databases cannot efficiently support at extremely high scale.
High-cardinality queries became slower, ingestion bursts caused backpressure, and sharding MongoDB clusters became expensive and operationally complex.
Countly needed:
- A storage engine purpose-built for analytical workloads
- A streaming backbone to decouple ingestion from processing
- A resilient architecture capable of handling peak bursts without dropping data
- Simpler operational overhead for customers
This is where Kafka and ClickHouse enter the story.
High-Level Architecture Overview
In the new design, Kafka becomes the central nervous system of Countly’s data flow. It manages event consistency, throughput, replayability, and pluggable consumers.
ClickHouse becomes the primary analytics engine, storing granular data in a columnar format for blazingly fast queries. MongoDB remains, but now as an operational metadata store and aggregation cache.
Everything in the architecture becomes more modular, more observable, and more scalable.
Components at a Glance:
- Ingestion Layer: Stateless pods performing validation, transformation, and efficient streaming.
- Kafka: High-throughput event backbone with replayable topic streams.
- Kafka Consumers: Compute units that build aggregations and store granular data.
- ClickHouse: Optimized for analytics, serving all drill and explorative queries.
- MongoDB: Stores application metadata, user profiles, and cached aggregations.
- API & Dashboard: Frontend and backend services reading from the combined data sources.
HA Deployment Blueprint:
- Kafka: 5-node cluster
- ClickHouse: 3–6 node replicated cluster
- MongoDB: 3-node replica set
- Countly Pods: Horizontally autoscaling
- Optional: Cross-region Kafka mirroring
Granular & Aggregated Data Reimagined
In the old architecture, aggregation and granular data competed within MongoDB. High-cardinality properties strained indexes. Drill queries performed full collection scans.
Countly’s new architecture now eliminates these limitations:
- Kafka is the single source of truth.
- ClickHouse stores all granular drill-level data - optimized for fast scans, grouping, and time-series analysis.
- MongoDB stores only cached pre-aggregations, significantly reducing load.
Aggregation consumers reprocess data when needed because Kafka offsets allow replay. This ensures strong consistency and resilience.
Deployment Type
Countly runs the same unified architecture across all deployments: MongoDB, Kafka, and ClickHouse working together as described above. What varies is how you run it:
Countly can be deployed on a single host with Docker Compose, which brings up the entire stack with minimal operational overhead. Suitable for low-volume applications, development environments, or customers who prefer a minimal stack.
For environments that need to scale across multiple nodes with high availability, Countly runs on Kubernetes using official Helm charts, where each component scales independently, and the databases are managed by their operators.
Both options deliver identical product features and dashboards; the difference is purely operational.
ClickHouse Storage Layout
ClickHouse uses schemas designed specifically for analytical workloads. The main table for granular analytics is stored under:
countly_drill.drill_events
It is optimized for fast scans, filtering, grouping, and time-series operations:
- MergeTree engine
- Monthly partitions
- Ordering by frequently filtered dimensions (app, event key, name, timestamp)
- JSON columns for dynamic segmentation and user properties
To further accelerate queries, the table includes:
- Bloom filter index on UID (for user-level queries)
- Min/max skip indexes on numerical and temporal fields
This layout ensures predictable and scalable performance even at very high data volumes.
Drill Events as the Granular Foundation of Analytics
drill_events represent the raw, event-level facts that power all explorative analytics in Countly.
Each drill event captures a complete analytical snapshot of what happened at the moment the event was recorded, regardless of whether the underlying engine is MongoDB or ClickHouse.
A drill event contains:
- Application and user identifiers
- Event key and name
- Timestamp of the event
- A flexible scalar field (
n) used by different features to store an identifier or a value(e.g., crash group ID for Crashes, widget ID for Surveys/NPS/Star Rating) - Segmentation data stored as structured JSON
- An expanded set of user-context fields, including selected custom properties captured at event time
This flexible design ensures that each event carries all information required for segmentation, filtering, grouping and user-level correlation without additional lookups. It also guarantees that analytics reflect the user’s exact state at the moment the event occurred, making drill events a more consistent and self-contained source of truth for high-cardinality analysis.
Query Compatibility Layer
To maintain full backward compatibility, Countly includes a compatibility layer that automatically translates MongoDB-style filters into their ClickHouse equivalents.
WhereClauseConverter
- Interprets operators such as
$in,$gt,$regex, and others - Produces type-safe, optimized ClickHouse
WHEREclauses - Ensures existing features behave consistently, regardless of whether MongoDB or ClickHouse is used behind the scenes
This allows customers to migrate to the high-volume architecture without changing dashboards, queries, or integrations.
Feature Coverage
Countly now clearly separates granular analytics and cached aggregations.
Granular, Explorative Analytics (ClickHouse-powered)
ClickHouse now directly powers these features:
- Drill
- Funnel
- Crashes
- Data Manager
- Compliance Hub
- Cohorts
- Views
- Session
- Ratings
- NPS & Surveys
- Formulas
- Retention
- Active Users
- Users → Events table
- Users → Sessions table
- Alerts
- Hooks
- Journey
Improvements
- Faster performance for large datasets
- No more 16MB MongoDB document limit
- More accurate and consistent metrics
- Better segmentation capability
- Consistent behavior across all explorative features
Hybrid Features
Features like User Profiles, Cohorts and Active Users combine granular reads with cached metadata for optimal responsiveness.
Restricted Features in the New Architecture
Data Governance Changes
The Data Manager remains a core part of Countly, but the new Kafka + ClickHouse pipeline imposes important restrictions on which operations can safely be supported.
These limitations exist because ClickHouse is designed as an immutable, append-optimized analytical store, and rewriting large volumes of historical data is not practical at scale.
When ClickHouse is enabled, the following actions are intentionally not supported, because they require rewriting historical events:
Modifying historical events or segments
(rename / merge / delete of already-stored values)
- Rewriting or reprocessing existing event properties
- Deleting or renaming past segment keys or values
- Transformation rules targeting existing data instead of incoming data
- Changing previously stored event records directly in the database (such as renaming historical events, updating old segment values, or rewriting past event properties...)
These actions are safe only in document databases like MongoDB, but do not align with ClickHouse’s immutable storage design.
Data Operations
Countly introduces a dedicated Mutation Manager to ensure that removing analytics data is performed safely, consistently, and without impacting system performance. Instead of deleting data immediately at the moment a request is made, Countly records the deletion intent in an internal queue and processes it through a controlled background job. This design prevents large or complex deletion operations from interfering with ongoing analytics, avoids partial or inconsistent results, and provides a predictable and auditable workflow.
Once a delete request is submitted, the Mutation Manager processes it in small batches, continuously validating progress, retrying failures with backoff, and marking the operation as complete only when it has finished successfully. If the system detects high load or conditions that could make the deletion unsafe to execute, the job automatically pauses and resumes later. This ensures that data removal never compromises system stability, responsiveness, or availability for end users.
The deletion job runs on a scheduled interval, configured to execute once per minute by default. Although this frequency can be adjusted through the Jobs interface, the default settings are recommended for most deployments, as they provide the safest balance between performance, throughput, and operational stability.
All deletion activity is observable through Countly’s monitoring endpoints, allowing administrators to track pending items, completed tasks, failures, and overall job health. By handling deletions as a governed background process rather than an instantaneous action, Countly offers a reliable, scalable, and compliance-ready approach to data removal suitable for production environments.
Performance & Scalability Characteristics
Countly’s internal deployment runs continuously at multi-billion-scale volumes.
Performance test results:
- 20K data points/second sustained ingestion (≈2B/day)
- Peaks up to 55K data points/second (≈140B/month)
- 100× – 1000× improvements across Drill, Funnels, Formulas, and User Timeline
- ClickHouse’s CPU-scaled query model allows predictable performance as the cluster grows
Operational improvements:
- No sharded MongoDB needed
- Kafka consumers scale horizontally
- Ingestion layer remains stateless
Typical Migration Journey
When customers migrate, the experience is designed to be safe, incremental, and predictable.
It usually begins with Countly deploying a parallel ingestion and storage pipeline alongside the customer’s existing production environment.
Initially, only new incoming data flows into Kafka. The customer continues using the legacy MongoDB-backed dashboards without interruption.
In the background, Countly launches a dedicated migration cluster, equipped with Debezium, which begins reading historical MongoDB collections and streaming them into Kafka for downstream processing.
This migration can move roughly 2TB of data per day without affecting live operations.
After a few days, the ClickHouse cluster has fully warmed up with both historical and real-time data.
At this point, Countly enables the new analytics engine, and the dashboards suddenly become faster, more responsive, and more capable.
Finally, once customer validation confirms parity of data, the legacy ingestion system is decommissioned. The customer now runs on the new high-throughput architecture without downtime or data loss.