Everything you need in place before starting the drill-events backfill on a Kubernetes deployment. Work through this article in order, then move on to running the migration.
If you have not read Migrating to Countly v26.01: Introduction yet, start there. It covers what the migration moves, how the service works, and how to choose a cutover scenario.
Before You Begin
- The v26.01 stack is installed and healthy. Follow Deploy Countly on Kubernetes with Helmfile first, including MongoDB via the Community Operator, ClickHouse via the Altinity Operator, and Kafka via Strimzi.
-
You have chosen a scenario. The three topologies in the introduction decide whether
LEDGER_CD_UPPER_BOUNDis set. Decide before you deploy the migration, not after. -
The ClickHouse
drill_eventstable exists. The migration service never creates it: theclickhouseplugin creates the schema on API startup. - Your cluster has headroom. The backfill is read-heavy on MongoDB and write-heavy on ClickHouse simultaneously. See Choose a Countly Helm Sizing Profile.
- You have disk space for both copies. During the migration the same events exist in MongoDB and ClickHouse. Budget the migrated data at roughly 10–20% of the MongoDB size after compression, plus transient headroom for the per-chunk staging tables, until you reclaim the MongoDB space at the end. The service's preflight measures actual free space on both sides and fails below 10%.
-
You can pull the migration image. Use
countly/countly-migration, or a tag you build yourself from the migration repository; add animagePullSecretsentry for a private registry.
Namespaces used throughout: mongodb and clickhouse for the backing services, and whichever namespace you deploy the migrator into. Adjust to your own layout.
Preparing the New Stack
Verifying the Target Stack and the ClickHouse Schema
Confirm the Countly, MongoDB, ClickHouse, and Kafka pods are running with kubectl get pods -A | grep -E 'countly|mongodb|clickhouse|kafka'. Then confirm the target table exists (substitute your ClickHouse pod name and credentials):
kubectl exec -n clickhouse <clickhouse-pod> -- \ clickhouse-client --password "$CLICKHOUSE_PASSWORD" \ --query "SHOW TABLES FROM countly_drill"
drill_events must be in the output. If it is not, check the Countly API pod logs for schema bootstrap errors before continuing. The service's own preflight fails on a missing table too, but confirming now saves a deploy cycle.
Setting Kafka Retention: For the Live Side, Not for the Backfill
The migration never touches Kafka. It reads documents out of MongoDB and inserts them straight into per-chunk ClickHouse staging tables, which it then attaches to the live table. No broker, no topic, and no connector is involved in the backfill path at any point.
This step is about the other half of your v26.01 stack. Live events from your SDKs do travel from the SDK through the API and Kafka into ClickHouse, and they keep arriving throughout the days the backfill runs. For those events, the Kafka log is the only replayable copy.
That matters in exactly one situation: the target ClickHouse is lost or has to be rebuilt mid-migration. Your data then still exists in two places: the history in the frozen source MongoDB, and the live events since cutover in the Kafka log. Recovery is to recreate the table, reset only the ClickHouse-sink connector's offsets to earliest, and re-run the migrator for the history. That path works only while the log still holds the window.
So raise drill-events retention to cover the whole migration before you start (the default is 14 days), and revert it at sign-off. If the backfill will run longer than retention, the gap is your exposure.
Replication factor is a separate, deliberate choice. RF ≥ 2 is recommended for large instances; if you run RF=1, record the accepted risk, because one broker disk loss forfeits the replay guarantee.
Pre-Copying the Stateful Set
The migration service moves drill events and nothing else. Everything the dashboard needs to interpret them has to be brought across separately, and most of it can be copied well before cutover with no user-facing impact: apps and app keys, app_users, event definitions, dashboard users, plugin configuration, and aggregated data.
At cutover you sync only the delta since this pre-copy, which is what keeps the ingestion pause down to minutes. Two ordering constraints apply at that point, both covered in Running the Migration: app_users must be complete before new ingestion starts, and aggregated data must land before new ingestion writes current-period documents.
Making the Source Data Reachable
Before the backfill can run, the migration has to be able to read your old drill_events* collections. How you arrange that is the single biggest planning decision in this migration, because it depends on how much data you have.
Start by separating two different needs, because they have different answers:
| What needs it | Which data | Access required | Typical share of total size |
|---|---|---|---|
| The migration service | countly_drill.drill_events* |
Read only, from anywhere it can reach over the network | The bulk of your data |
| Countly v26.01 itself |
countly (apps, members, aggregates) and countly_drill's drill_meta* / drill_bookmarks
|
Read/write, in the cluster's own operational MongoDB | The smaller share |
That split is what makes the large-data approaches possible: the events never have to be copied into the new cluster at all, only read out of wherever they already live.
One more thing shapes every option below: the migration uses one MongoDB connection string for three different things.
| Setting | Default | What the migration does with it |
|---|---|---|
MONGO_DB |
countly_drill |
Reads the drill_events* collections: the actual source data |
MONGO_COUNTLY_DB |
countly |
Reads apps and events to resolve per-event collection hashes during transform |
MANIFEST_DB |
countly_drill |
Writes the ledger and DLQ (mig_ranges, mig_dlq_docs, mig_run_config, and mig_collection_est) |
The second of these is easily overlooked, and it is not optional. Old drill documents identify their app and event with a 40-character SHA1 hash. The migration reconstructs the real values by reading every app from countly.apps and every custom event name from countly.events, hashing eventName + appId, and building a lookup table at startup. If that database is missing or empty on the connection you give it, the lookup table is empty and hashed events cannot be resolved.
So whichever option you pick, the connection string you hand the migration must expose both countly_drill and a matching countly, and it must accept writes for MANIFEST_DB.
Option A is the right starting point for most deployments. Options B and C exist for cases where copying the event collections is impractical or unnecessary.
Option A: Restoring From a Volume Snapshot (Recommended)
The default choice. Only the drill events move to ClickHouse. v26.01 still runs on MongoDB for apps, members, aggregates, drill metadata, and bookmarks, so that data has to exist in the new cluster regardless. A volume clone brings countly and countly_drill across together in one operation, which is exactly what both Countly and the migration need.
This is a block-level copy, so it moves terabytes in the time your storage layer takes to clone a volume, rather than the time mongorestore takes to re-insert every document and rebuild every index.
- Take a snapshot of the source MongoDB data volume with your CSI driver or cloud provider. For a replica set, snapshot one secondary, after stopping writes to it or using a filesystem-consistent snapshot.
-
Create a PVC in the new cluster from that snapshot:
apiVersion: v1 kind: PersistentVolumeClaim metadata: name: mongodb-restored-0 namespace: mongodb spec: storageClassName: <your-storage-class> dataSource: name: <your-volume-snapshot> kind: VolumeSnapshot apiGroup: snapshot.storage.k8s.io accessModes: ["ReadWriteOnce"] resources: requests: storage: 2Ti - Bring MongoDB up on that PVC and let the operator reconcile.
Two issues commonly arise here:
-
Replica set identity travels with the volume. The
localdatabase on the cloned disk still contains the old replica set name and member hostnames. Starting a new member on it without reconciling that gives you a node that will not join. Either match the old replica set name in your MongoDB resource, or start the node standalone, drop thelocaldatabase, and re-initiate the replica set. - MongoDB version compatibility. Data files can be opened by the same major version or upgraded one major at a time. Bring the clone up on the version that wrote it, confirm it is healthy, and only then upgrade.
If you snapshot inside the ingestion pause (pause old ingestion, snapshot, then resume ingestion on the new stack), the copy can never hold a mirror copy of a natively ingested event. That is the cleanest possible run: no bound, no duplicates, and nothing to deduplicate later. A snapshot taken after ingestion resumed has a duplicated tail and needs the dedupe pass described in Running the Migration.
Note that a frozen copy changes what the audits compare against: the tool sees the copy, not the live old cluster, so zero counts after the snapshot moment mean "snapshot taken here", not lost data. Scope any comparison against the live old deployment to cd below the snapshot moment.
Option B: Pointing the Migration at Your Existing MongoDB (No Copy)
For very large event volumes, in combination with Option A or C. Nothing is copied for the events themselves: the migration streams them straight out of the old cluster into ClickHouse. This works well because the old cluster already has countly_drill and countly side by side, exactly as the migration expects.
MONGO_URI: "mongodb://migrator:PASS@old-mongo-1:27017,old-mongo-2:27017/admin?replicaSet=rs0&ssl=false" MONGO_DB: "countly_drill" MONGO_COUNTLY_DB: "countly" MANIFEST_DB: "countly_migration_manifest"
Requirements and caveats:
-
The migration needs write access to the old cluster. There is only one connection string, so the ledger and DLQ are written through the same URI, into
MANIFEST_DBon the old cluster. Grant the migration user read oncountly_drillandcountly, and write on that one database. Override the default here.MANIFEST_DBdefaults tocountly_drill, which would write migration state into the old production drill database; a distinct name likecountly_migration_manifestkeeps it identifiable and safe to drop later. - Network reachability. The migration pods must reach every member the replica set advertises on port 27017, not just the seed hosts, because the driver connects using the hostnames in the replica set config.
-
Read preference is automatic. On a replica set the engine selects
secondaryPreferredby itself. Because the source is frozen after cutover, secondary reads are exact. Only setMONGO_READ_PREFERENCEto override that deliberately; preflight warns if you have forcedprimary. -
Read load on a live cluster. Reading a few thousand documents per second from a secondary is real load, and the app and event lookup adds a full scan of
countly.appsat startup. Prefer a dedicated hidden or analytics secondary if your replica set has one. -
The new cluster still needs its own
countly. Option B only solves reading the events. Countly v26.01 serves the dashboard from its own operational MongoDB, socountly(apps, members, aggregates) andcountly_drill'sdrill_meta*/drill_bookmarksstill have to exist in the new cluster. Those are the smaller share of your data, so bring them across with Option A or Option C. That combination (operational databases cloned or restored, events read in place) is the usual reason to choose this option.
Option C: Dumping and Restoring
For small deployments, and for the operational databases alongside Option A or B. A dev or staging environment, a single small app, or a total dataset in the tens of gigabytes restores fine this way. Beyond that, mongorestore becomes impractical for the events: it re-inserts every document and rebuilds every index, which takes far longer than the equivalent volume clone and needs scratch space for the archive on both ends.
# On the OLD server mongodump --db countly --archive=/backup/countly.archive mongodump --db countly_drill --archive=/backup/drill.archive # Into the MongoDB pod kubectl exec -i -n mongodb <mongodb-pod> -- \ mongorestore --uri "mongodb://app:$MONGODB_PASSWORD@localhost:27017/admin" \ --archive --drop < /backup/countly.archive
See the mongorestore documentation for the details. If you use this route at any real size, run it as a Kubernetes Job with the archive on a PVC rather than streaming through kubectl exec, which will not survive a dropped connection.
If you are pairing this with Option B, you only need the countly archive plus the drill_meta* and drill_bookmarks collections, not the drill_events* collections, which the migration reads in place.
Building the Read Index
The migration pages through each collection on a {cd: 1, _id: 1} compound index. It builds this index itself for any collection that lacks one, and it will not start processing a collection until that collection's index is ready. On a large deployment that means a fresh run can spend its first hours (or, at 10 TB, its first one to three days) building indexes with no rows moving, which looks exactly like a stuck migration.
Start these builds during the preparation phase, not after cutover. Build them in the background, throttled, on secondaries where possible:
kubectl exec -i -n mongodb <mongodb-pod> -- mongosh countly_drill --quiet <<'EOF'
db.getCollectionNames().filter(n => n.startsWith("drill_events")).forEach(n => {
print("indexing " + n);
db[n].createIndex({ cd: 1, _id: 1 });
});
EOFThe service can also start the builds for you once it is running (Build indexes on the dashboard, or POST /control/build-indexes, with progress at /api/index-progress), but starting early overlaps the wait with work you were doing anyway.
No collection consolidation is ever needed. The service discovers every collection matching the drill_events prefix and maps chunks across all of them.
Giving MongoDB and ClickHouse Enough Memory and CPU
The backfill loads both databases at once: MongoDB serves a continuous paged read while ClickHouse absorbs a sustained insert stream, holds a staging table per in-flight chunk, and merges the parts it creates. MongoDB reads are usually the first bottleneck, but ClickHouse needs headroom too. An under-resourced ClickHouse spends the migration in backpressure, pausing the inserts while it catches up on merges.
What MongoDB is short of is almost always memory, not CPU. The scan runs over collections far larger than RAM, so throughput is governed by how much of the index and working set stays resident in the WiredTiger cache. Starve the cache and every page turn becomes a disk read, which no number of cores compensates for.
Before starting, raise both databases for the duration of the migration:
- MongoDB: raise the WiredTiger cache first, then the pod's memory request and limit. Keep the limit at least twice the cache, because WiredTiger needs headroom beyond it for connections, sorts, and page reconstruction, and a limit set too close produces an OOMKill rather than a faster scan. CPU matters least of the three.
-
ClickHouse: raise CPU and memory in your
ClickHouseInstallation.
Then confirm the pods actually restarted with the new values:
kubectl get pods -n mongodb -o jsonpath='{range .items[*]}{.metadata.name}{"\t"}{.spec.containers[*].resources.limits}{"\n"}{end}'Scale both back to your steady-state sizing profile after the migration completes.
Preparing the Migration Manifests
The migration service ships ready-to-apply manifests. It is deployed on its own rather than as part of the Countly Helm release, and needs no chart and no operator: pods coordinate through chunk leases in MongoDB, so a plain Deployment is all the orchestration required.
git clone https://github.com/Countly/migration.git cd migration
| Manifest | Shape | Use it when |
|---|---|---|
k8s/migration.yaml |
ConfigMap + Secret + Deployment + Service | Recommended. Pods keep serving the dashboard after completion, so you can verify and sign off in the UI |
k8s/job.yaml |
Job, using EXIT_ON_COMPLETE
|
Fire-and-forget automation. Requires the ConfigMap and Secret from migration.yaml. The dashboard disappears with the pods when the Job finishes |
Edit the ConfigMap and Secret before applying. The Secret carries the two required values:
apiVersion: v1 kind: Secret metadata: name: drill-migrator-secrets type: Opaque stringData: MONGO_URI: "mongodb://mongo.example.internal:27017" CLICKHOUSE_URL: "http://clickhouse.example.internal:8123" #CLICKHOUSE_USERNAME: "default" #CLICKHOUSE_PASSWORD: ""
The ConfigMap carries everything else. The settings you are most likely to touch:
| Variable | Default | What it controls |
|---|---|---|
LEDGER_RUN_ID |
ledger-v1 |
The resume key. Keep it identical across every pod and every restart of the same migration; changing it starts a separate run |
MANIFEST_DB |
countly_drill |
Where the ledger and DLQ live. Override when reading from a live old cluster (Option B) |
LEDGER_CD_UPPER_BOUND |
unset | Set only in scenario 2 (mirrored cutover). Epoch ms or ISO. Documents at or after it are never migrated |
LEDGER_UNBOUNDED_OK |
false |
Declares up front that nothing mirrors traffic, so the startup guard does not hold the run |
LEDGER_START_PAUSED |
false |
Deploy now, start later: pods serve the dashboard and wait until Start is pressed once for the whole fleet. Set this: it is how you rehearse without starting a migration |
POD_ID |
pod name | Defaults to the hostname, which in Kubernetes is the unique pod name; nothing to configure |
EXIT_ON_COMPLETE |
false |
Set by job.yaml. Leave unset for the Deployment, whose pods should stay up for verification |
MONGO_READ_PREFERENCE |
auto | Leave unset. The engine picks secondaryPreferred on a replica set by itself |
LEDGER_CHUNK_DOCS_TARGET |
2000000 |
Target documents per chunk |
LEDGER_MAX_CHUNK_DAYS |
7 |
Upper bound on a chunk's time span; guards against a bad estimatedDocumentCount producing one enormous chunk |
MONGO_PAGE_SIZE |
10000 |
Documents per read page |
LEDGER_INSERT_INFLIGHT |
3 |
Concurrent inserts into the staging table |
DRY_RUN / DRY_RUN_SAMPLE_PCT
|
false / 2
|
Sampled rehearsal against a throwaway target; nothing is stored |
.env.example in the repository is the full commented reference. Resuming needs no setting of its own: it is what LEDGER_RUN_ID does.
On resources: the shipped Deployment requests 2 CPU / 2 Gi and limits 4 CPU / 6 Gi per pod. The image caps the Node heap at 4 GB, which is why the memory limit sits above it: setting it at or below 4 Gi guarantees OOMKills under load. Raising a single pod's limits past this buys little; throughput comes from more pods, which the next section covers.
Deciding How Many Pods to Deploy
The migrator is horizontally scalable, and how many pods you run is a decision worth making before you apply anything rather than discovering mid-backfill. Adding pods later is safe and takes one command, but sizing the nodes for them is not something you want to do under time pressure.
How the Pods Divide the Work
Every collection is mapped into cd-bounded chunks upfront, and pods claim them from that one shared grid. A claim is an atomic, leased write in MongoDB, so two pods can never take the same chunk, and a pod that dies has its lease expire so another picks the chunk up and redoes it from scratch.
There is nothing to partition by hand. Pods claim in collection order, newest data first within each collection, and spill into the next collection the moment the current one has nothing claimable. A dataset made of hundreds of small collections parallelizes exactly as well as one enormous one.
What Every Pod Must Share
| Setting | Requirement |
|---|---|
LEDGER_RUN_ID |
Identical on every pod. It is the key that binds them to one run; a pod with a different value quietly starts a second, separate migration |
MONGO_URI, MANIFEST_DB
|
Identical: this is where the shared chunk grid lives |
CLICKHOUSE_URL, CLICKHOUSE_DB, CLICKHOUSE_TABLE
|
Identical |
LEDGER_CD_UPPER_BOUND |
Identical. A pod missing the bound migrates past the boundary while its neighbors do not, which is why you check the bounded · cd < … badge on every pod, not just one |
POD_ID |
Leave it unset. It defaults to the hostname, which in Kubernetes is the unique pod name |
Sharing one ConfigMap and one Secret across the Deployment, as k8s/migration.yaml does, satisfies all of this by construction. That is the reason to scale by replica count rather than by hand-rolling a second Deployment with its own config.
Setting the Replica Count
Set the count in the Deployment:
apiVersion: apps/v1 kind: Deployment metadata: name: drill-migrator spec: replicas: 3 # …
Or, for the Job variant, keep the two values equal. Every pod exits 0 on completion, so the Job finishes only when all of them do:
spec: parallelism: 3 completions: 3
You can change the count at any time, including mid-run, with kubectl scale deployment/drill-migrator --replicas=5.
Scaling down is equally safe. A terminated pod's in-flight chunk is redone by whoever takes the lease next, so there is nothing to drain and no graceful-shutdown dance to get right.
Spreading Them Across Nodes
Pods scale across nodes, not within one. A single pod saturates roughly four cores on BSON decode, so two pods on the same node mostly contend for the same CPU and buy you very little. The shipped manifest requests 2 CPU and limits 4, which lets the scheduler make that judgment for you on a reasonably sized cluster, but on a large node it will happily pack several pods together.
If you see that happen, spread them explicitly:
spec:
template:
spec:
topologySpreadConstraints:
- maxSkew: 1
topologyKey: kubernetes.io/hostname
whenUnsatisfiable: ScheduleAnyway
labelSelector:
matchLabels: { app: drill-migrator }Confirm the spread with kubectl get pods -l app=drill-migrator -o wide.
How Many Is Enough
More pods stop helping once something else becomes the constraint, and on this workload that is almost always MongoDB read capacity first and ClickHouse merge throughput second. Pods are cheap; the databases behind them are not. Two consequences worth planning around:
- If MongoDB is already saturated, adding pods makes the backfill no faster and the source cluster no happier. Fix the WiredTiger cache before adding replicas.
- If ClickHouse starts spending its time in backpressure (inserts pausing while merges catch up), you have passed the useful point. Fewer pods and a healthier target finish sooner than more pods and a permanently throttled one.
A practical approach: start with one pod, watch the rate on the dashboard once the run is going, and add replicas while the rate still climbs roughly in proportion. When it stops doing that, you have found the ceiling, and it is a database ceiling rather than a pod-count one.
Whatever count you settle on, every pod's dashboard shows and controls the entire run, because the state is in MongoDB rather than in any one pod. There is no leader to find and no per-pod view to reconcile.
Deploying Without Starting the Run
Apply the manifests with the start gate closed, so the pods serve the dashboard and answer their API without touching your data. Set LEDGER_START_PAUSED: "true" in the ConfigMap, then apply and watch:
kubectl apply -f k8s/migration.yaml kubectl logs -f deploy/drill-migrator
The log line confirming the gate is holding:
LEDGER_START_PAUSED: holding before mapping — press Start on the dashboard (POST /control/resume) to begin
The hold is before mapping, deliberately. Mapping is what builds the {cd, _id} index on the source and cuts the chunk grid, and a run nobody has started yet should be doing neither. In this state the pods have read nothing, written nothing, and indexed nothing; they report pauseReason: not-started and wait.
What they will do while held is serve the dashboard and answer three things: preflight, index builds, and the dry run. Those are exactly what you need before cutover, and all three are safe to run against a live old deployment.
Reach the dashboard, which any pod serves in full:
kubectl port-forward svc/drill-migrator 8080:8080 # http://localhost:8080
The gate lives in the ledger rather than in the pod, so a pod that restarts while held stays held, pods that join later hold too, and one Start covers the whole fleet. Nothing begins until you press Start; that step belongs to Running the Migration.
Running the Built-In Preflight
With the pods up and holding, run preflight from the dashboard's Migration Guide tab, or through the port-forward:
curl -s localhost:8080/api/preflight | jq '.checks[] | {label, status, detail}'It is read-only, so run it as often as you like. It reports pass, warn, or fail for each of:
| Check | What a failure means |
|---|---|
| MongoDB source reachable | Wrong URI, credentials, or network policy; also reports how many drill_events* collections it found |
{cd,_id} index on all collections |
Warn only: the service will build the missing ones, but you will wait |
| Estimated documents to migrate | Informational; the estimate can be off after an unclean mongod shutdown, which is why chunk sizing is span-guarded |
| Source frozen & clocks sane |
Fail. The newest source cd is within 60 s of ClickHouse server time, so the boundary between migrated and live data is not trustworthy: old ingestion is still running, or the clocks are skewed |
| Old ingestion stopped | Fail. A four-second probe saw a collection still growing. In bounded (mirror) mode this check is replaced by a bound-sanity check, because the source is expected to keep growing |
| New ingestion flowing into ClickHouse | Warn only: either the SDK flip has not happened yet or traffic is genuinely zero |
Documents without cd
|
Warn: a dedicated sweep chunk migrates them, strictly after all regular chunks |
| Replica set detected | Warn if you forced MONGO_READ_PREFERENCE=primary; remove it and let the engine pick secondaryPreferred
|
| MongoDB / ClickHouse disk headroom | Fail below 10% free, warn below 20%. The most preventable mid-migration incident there is |
| ClickHouse target table |
Fail. drill_events does not exist; start the new stack first |
| Insert-dedup canary | Warn if dedup tokens are inert on this target. Safe either way, because chunk redo covers it |
| Dry run | Warn until you have run one |
Rehearsing With a Dry Run
A dry run maps a sampled subset of the work and takes it through the full transform and ClickHouse validation path against a Null-engine clone of the target, storing nothing. It is the last prerequisite, and it is safe against a live old deployment.
Trigger it on the deployment you already have up with the dashboard's Dry run button, or:
curl -s -X POST localhost:8080/control/dry-run -H 'content-type: application/json' -d '{}'
curl -s localhost:8080/api/dryrun | jq # poll until status is "completed"You do not need to set DRY_RUN in the ConfigMap, and you do not need a second Deployment or Job. The endpoint opens its own source and target connections so it can never disturb the main run's state. The rehearsal is recorded under its own run id (<LEDGER_RUN_ID>-dry) with its own start gate, so rehearsing can never silently authorize the real run.
Two conditions apply. The main run must not be actively copying (holding at the gate is fine, and so is paused), and only one dry run may be in flight. If either is violated, the call returns {"started": false, "reason": "…"} rather than queueing.
To change the sample size, set DRY_RUN_SAMPLE_PCT in the ConfigMap before the pods start. It accepts 0.1 to 5 and defaults to 2.
Then review the report with whoever owns sign-off, using curl -s localhost:8080/report | jq.
It lists skip reasons, per-key coercions, and a DLQ summary. What you are looking for is anything systematic (a whole event type coercing, or a skip reason with a large count), because that is far cheaper to fix now than mid-backfill.
Re-run preflight afterwards: its Dry run check turns green once the sampled chunks are done, which is the signal that the prerequisites are complete.
Next Steps
With the prerequisites met, continue to Running the Migration.