Kubernetes: Migration Prerequisites

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_BOUND is set. Decide before you deploy the migration, not after.
  • The ClickHouse drill_events table exists. The migration service never creates it: the clickhouse plugin 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 an imagePullSecrets entry 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.

  1. 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.
  2. 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
  3. Bring MongoDB up on that PVC and let the operator reconcile.

Two issues commonly arise here:

  • Replica set identity travels with the volume. The local database 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 the local database, 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_DB on the old cluster. Grant the migration user read on countly_drill and countly, and write on that one database. Override the default here. MANIFEST_DB defaults to countly_drill, which would write migration state into the old production drill database; a distinct name like countly_migration_manifest keeps 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 secondaryPreferred by itself. Because the source is frozen after cutover, secondary reads are exact. Only set MONGO_READ_PREFERENCE to override that deliberately; preflight warns if you have forced primary.
  • 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.apps at 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, so countly (apps, members, aggregates) and countly_drill's drill_meta* / drill_bookmarks still 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 });
});
EOF

The 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.

Was this page helpful?
Reach out to us for any other questions.
Helpful?

Looking For More Help?