Bulk Batch Replication is a batch (polling) sibling of the streaming Realtime Export and Delta Lake Export paths. Instead of forwarding one resource change at a time through a broker channel, a clustered scheduled trigger periodically replicates a whole window of HFJ_RES_VER (resource version) rows in bulk to a configured target repository.
Where the streaming path optimizes for low latency (each change is written as it happens), Bulk Batch Replication optimizes for write throughput: it coalesces many resource changes into large FHIR transaction Bundles and hands each Bundle to a single repository call, so the target receives few large writes rather than many small ones. This suits targets where large, infrequent writes are preferable to a steady stream of single-row operations.
The replication target is any implementation of the HAPI FHIR ca.uhn.fhir.repository.IRepository interface, selected by configuration. Any IRepository works — an analytical data lake (for example, an Iceberg- or Delta-backed repository) is one motivating use case, but the module makes no assumption about what the target does with the Bundles. Out of the box the target defaults to a no-op repository that discards everything, so the module is inert until a real target is configured.
Bulk Batch Replication is configured as its own module type (REALTIME_EXPORT_BULK_BATCH_REPLICATION), independently enabled per deployment. It depends on a persistence module (the source of the HFJ_RES_VER rows) and on the Cluster Manager module (for clustered scheduling).
The persistence module must be a relational one — extraction reads the HFJ_RES_VER table directly, which MongoDB persistence does not have. Wiring the module against a MongoDB persistence module fails the module's startup with a configuration error.
See the Bulk Batch Replication module reference for the module's configuration categories and dependencies.
A cluster-singleton scheduled trigger fires on a fixed interval. Each fire starts one gated Batch2 job, so replication runs benefit from Batch2's retry, gating, and observability; the trigger only schedules, it does not do the extraction itself.
Scheduled trigger (every schedule_interval_ms, on one node in the cluster)
→ start one Batch2 job
├ fan out: one unit of work per partition (or shard)
├ per partition: drain the resource-version rows written since the last run,
│ in pages; each page → one FHIR transaction Bundle → target repository;
│ record durable progress (a watermark) after every written page
└ report: rows written, partitions, recovered errors
Each run reads directly from the resource version history (HFJ_RES_VER) using a resumable cursor that is exactly-once at each page boundary: the next run continues precisely where the previous one stopped, never skipping or repeating a row, even when many versions are written in the same instant under concurrent load.
Because extraction is against the version history table (not current state), every version of a resource is replicated, not just its latest state. The lake accumulates full history, which suits analytical, append-oriented tables.
Each page of version rows is assembled into one FHIR transaction Bundle and submitted in a single repository call, so the target batches its writes rather than applying changes one row at a time. Live versions become versioned updates and deleted versions become deletes, so both updates and deletes propagate: for each resource, the target reflects the source's full history, ending at its current state.
The job replicates every partition independently and in parallel (per shard under MegaScale); a failure in one partition is reported without failing the others in the same run.
partition_scopepartition_scope restricts replication to specific partition ids; leaving it empty replicates all partitions. It requires a partitioned persistence module — configuring a scope without partitioning enabled fails module startup, since such a scope could never match any rows.
A partial scope (some but not all tenants of a shard) is more expensive to drain than a whole-shard scope, because it filters on a column the database does not index by partition. The cost is bounded and matters mainly during the initial backfill (whose window is the entire table), not in steady state. For a large existing deployment, prefer backfilling a partial scope deliberately — leave the schedule disabled and drive it with $sdh.bulk-batch-replication-start, then enable the schedule once the watermark has caught up. Adding a partition to the scope later re-replicates that partition's history (harmless for idempotent targets, but expect a backfill-sized run).
The set of shards (and each shard's co-located partitions) is snapshotted when a run starts, and a whole-shard drain records its watermark against that shard's set of co-located partition ids. As a result, changing the MegaScale topology while replication is running is not supported:
Task is left behind unused.Apply topology changes with replication quiesced: disable the schedule (enabled=false), confirm no run is in flight, make the change, then re-enable. Expect the next run to be backfill-sized for any shard whose partition membership changed — size stale_instance_threshold_ms accordingly, or drive the catch-up manually with $sdh.bulk-batch-replication-start before re-enabling the schedule.
Replication runs never pile up: a new run will not start while a prior one is still in flight (this applies to both the schedule and the manual start operation). A run left behind by a crashed or restarted node is automatically treated as orphaned after stale_instance_threshold_ms and cancelled so a fresh run can start — without losing progress, since each written page is durably recorded and the replacement resumes from there.
Because the cluster-wide guard is best-effort, avoid issuing a manual start around a scheduled fire (or disable the schedule while running a manual backfill): two runs briefly overlapping cannot corrupt the target, but they redundantly re-replicate the same rows.
Because each run executes as a Batch2 job, its progress, status, and history are visible through the standard Batch2 job instance surface. In addition, the module emits the following OpenTelemetry metrics, each tagged with smilecdr.module_id:
| Metric | Type | Description |
|---|---|---|
smilecdr.bulk_batch_replication.rows_replicated | counter | Resource-version rows replicated to the target repository. |
smilecdr.bulk_batch_replication.partition_drain.duration | histogram (seconds) | Time taken to drain a single partition. |
smilecdr.bulk_batch_replication.partition_recovered_errors | counter | Partitions that recovered from an error (the next run resumes them from their last durable watermark). |
smilecdr.bulk_batch_replication.job.duration | histogram (seconds) | End-to-end replication job duration. |
Per-partition progress, drain summaries, and final job status are also written to the Realtime Export troubleshooting log.
Besides the recurring schedule, a run can be started on demand with the FHIR operation $sdh.bulk-batch-replication-start at the server (system) level. This is useful for ad-hoc backfills and for QA. The operation enforces the same overlap guard as the scheduler — it fails with a precondition error if a prior instance is still in flight.
When more than one bulk-batch-replication module is configured, each is an independent replication target. Identify the one to run with the target parameter, set to that module's id. The parameter may be omitted when exactly one module is configured; if several are configured and target is omitted, the operation fails and lists the configured targets. An unknown target fails with a not-found error.
| Parameter | Type | Description |
|---|---|---|
target | string | The module.id of the bulk-batch-replication instance to run. Optional when exactly one is configured; required to disambiguate when several run concurrently. |
endBound | dateTime | Rows with RES_UPDATED at or after this instant are excluded from the run. Defaults to now. |
pageSize | integer | Maximum number of HFJ_RES_VER rows fetched per page, per partition. Defaults to the configured page size. |
The enabled setting gates only the recurring scheduled trigger. The module and the manual start operation remain available even when the schedule is disabled, so an operator can keep ad-hoc triggering while leaving the recurring schedule off.
All settings are backend-agnostic and apply regardless of the repository target.
module.bulk_batch_replication.config.enabled=false
module.bulk_batch_replication.config.schedule_interval_ms=900000
module.bulk_batch_replication.config.stale_instance_threshold_ms=2700000
module.bulk_batch_replication.config.trailing_lag_slop_ms=300000
module.bulk_batch_replication.config.page_size=500
module.bulk_batch_replication.config.partition_scope=
module.bulk_batch_replication.config.target_repository=ca.cdr.rte.repository.NoOpRepository
| Property | Description | Default |
|---|---|---|
enabled | If enabled, the clustered scheduled trigger periodically fires replication runs. Gates only the schedule, not the module or the manual start operation. | false |
schedule_interval_ms | How often the trigger fires a new replication run; also your replication latency floor. See Choosing the duration values. | 900000 (15 min) |
stale_instance_threshold_ms | How long a run may look stalled before it is assumed dead (e.g. a crashed node) and cancelled so a fresh run can start. See Choosing the duration values. | 2700000 (45 min) |
trailing_lag_slop_ms | How far back from "now" each run holds off, so a write that has not committed yet is never skipped; deferred rows are picked up by the next run. See Choosing the duration values. | 300000 (5 min) |
page_size | Maximum number of resource-version rows fetched per page, per partition. Governs per-run memory and checkpoint granularity, not the target's file/object size (the target decides how it batches Bundles into files). | 500 |
partition_scope | Comma-separated partition ids to replicate. Leave empty to replicate all partitions. | (all partitions) |
target_repository | The replication target: either a fhir-repository: URL or a fully-qualified IRepository class name (see Repository target). Defaults to a placeholder that discards everything. | ca.cdr.rte.repository.NoOpRepository |
Three settings are durations. Pick each from your deployment's characteristics rather than accepting the defaults blindly.
schedule_interval_ms — your replication latency floor. A newly written row first appears in the target after at most one interval plus trailing_lag_slop_ms. Shorter intervals give fresher data and smaller per-run windows but run more often; longer intervals give fewer, larger runs. Keep the 15-minute default for analytical targets; drop to a few minutes only if you genuinely need fresher data and the target can absorb more frequent writes. There is no benefit to setting it below how long a typical run takes — runs never overlap (a fire is skipped while a prior run is still going), so an over-short interval just yields back-to-back runs.
trailing_lag_slop_ms — how far back each run stays clear of "now". It only has to cover the brief moment between a row being stamped and its database transaction committing, plus any clock skew between the application and database servers — seconds in a healthy, time-synchronized deployment. It does not need to cover the duration of your longest transaction. The 5-minute default is deliberately conservative; if you want lower latency and your app and database clocks are synchronized (e.g. NTP), lowering it to 30–60 s is safe. Raise it only if you observe recently-written rows intermittently missing from a run and appearing in the next one — a sign the real gap is larger than configured (for example, unsynchronized clocks).
stale_instance_threshold_ms — how long a run may look stalled before it is assumed dead and cancelled. Set it above your longest legitimate run, which is almost always the initial backfill rather than a steady-state run. Too low, and a healthy long backfill gets cancelled and restarted; too high, and recovery from a genuinely crashed run is delayed by that much. Size it from an observed or estimated worst-case backfill duration with generous headroom — the 45-minute default suits moderate datasets; increase it for very large initial loads.
The target is resolved from target_repository, which accepts either form:
fhir-repository: URL, resolved by HAPI's repository loader SPI — for example a filesystem or in-memory target (fhir-repository:exp-kalm-filesystem:/path, fhir-repository:memory:name). Because replication emits every resource version as a separate transaction-Bundle entry, the target must accept a bundle carrying multiple versions of the same resource; lake-style (append/filesystem) loaders do, so those are the appropriate URL targets.ca.uhn.fhir.repository.IRepository, constructable via the standard customer-bean mechanism (a no-argument or autowired constructor; a FhirContext is offered as an available dependency).The default NoOpRepository accepts and discards every Bundle, which keeps the module inert until a real target repository is configured.
$expunge does not propagate. A physical expunge removes rows from HFJ_RES_VER, so there is no version row left for the cursor to observe. Expunged resources are not deleted from the target. This is the same limitation as the streaming path.