BulkLoad (Dev)
Overview
The BulkLoad feature works in conjunction with BulkDump to provide a complete data migration solution. BulkLoad takes manifest files and SST files generated by BulkDump and efficiently loads them into a target FoundationDB cluster.
Input and Requirements
When a user wants to start a bulkload job, the user provides:
Job ID: The unique identifier from the corresponding BulkDump job
Key Range: The range of keys to load (must be within the dumped range)
Source Path: Either a local directory or blobstore URL containing the dump files
Required Configuration: BulkLoad requires the following server knobs to be enabled:
--knob_shard_encode_location_metadata=1: Enables shard-aware location metadata--knob_enable_read_lock_on_range=1: Enables exclusive range locking during load operations
Input File Structure: BulkLoad expects the input files to be organized as produced by BulkDump.
How to use?
Currently, FDBCLI tools and low-level ManagementAPIs are provided to submit a job or clear a job. These operations are achieved by issuing transactions to update the bulkload metadata and taking exclusive locks on the target range. Submitting a job involves validating the input parameters, taking an exclusive read lock on the target range, and writing job metadata. When submitting a job, the API checks if there is any ongoing bulkload job or conflicting locks. If yes, it will reject the job. Otherwise, it accepts the job. Clearing a job releases the range lock and marks the job as cancelled in the metadata.
FDBCLI provides following interfaces to do the operations:
Submit a job: bulkload load (JobID) (BeginKey) (EndKey) (RootFolder) // …where JobID is from BulkDump, and RootFolder is a local directory or blobstore URL
Clear a job: bulkload cancel (JobID)
Enable the feature: bulkload mode on | off // “bulkload mode” command prints the current value (on or off) of the mode
Check status: bulkload status // Shows current running job information
View history: bulkload history // Shows completed job history
For detailed usage examples and quickstart guide, see BulkLoad User Guide.
ManagementAPI provides following interfaces to do the operations:
Submit a job:
submitBulkLoadJob(BulkLoadJobState jobState)Clear a job:
cancelBulkLoadJob(UID jobId)Enable the feature:
setBulkLoadMode(int mode)// Set mode = 1 to enable; Set mode = 0 to disableGet job status:
getBulkLoadJobStatus(Database cx)BulkLoad job metadata is generated by
createBulkLoadJob()
Mechanisms
Workflow
Users submit a BulkLoad job via
submitBulkLoadJob()specifying the source JobID, target range, and data locationThe API validates the job parameters and checks for conflicting BulkLoad/BulkDump jobs
An exclusive read lock is taken on the entire target range using
takeExclusiveReadLockOnRange()Job metadata is persisted to bulkload job space (
\\xff/bulkLoadJob/prefix) and task space is initializedDD’s
bulkLoadJobManager()detects the new job and downloads the global job-manifest.txt fileDD parses the job manifest to build a map of manifest entries by key range
DD creates BulkLoad tasks by grouping manifest entries (up to
MANIFEST_COUNT_MAX_PER_BULKLOAD_TASKper task)Each task is persisted to task metadata space (
\\xff/bulkLoadTask/prefix) and triggers data movementDD’s
doBulkLoadTask()coordinates with data movement system to load SST files into target shardsStorage servers receive data movement requests containing BulkLoad task information
Storage servers download SST files, validate integrity, and apply data using storage engine ingestion
Tasks complete and are marked as
BulkLoadPhase::CompleteorBulkLoadPhase::ErrorWhen all tasks finish, the job is finalized and moved to job history, and the range lock is released
Per-Task Range Clear
Before ingesting SST files for a task, the storage server clears the target
shard range and waits for the clear to durably commit. This ensures that
stale data is not retained when the new SST files are ingested over an
existing range. The clear-then-ingest sequence is performed per shard inside
the SS (see storageserver.cpp):
getKeyValueStore()->clear(keys)— clear the target rangeWait on
durableVersionso the clear is committed before ingestcompactRange(keys)to optimize storage before bulk insertingestSSTFiles(localBulkLoadFileSets)— ingest the new data
Because the SS clears each shard internally, BulkLoad does not require the
destination cluster or range to be empty when the job starts. The exclusive
range lock taken during submitBulkLoadJob() (see below) ensures no
concurrent writes can race with the clear-and-ingest sequence. As a result,
fdbrestore start --mode bulkload is permitted to run against a
non-empty destination database — unlike rangefile restore, which requires
an empty destination.
Range Locking
BulkLoad uses FoundationDB’s range locking mechanism to ensure data consistency:
registerRangeLockOwner()registers the BulkLoad system as a lock owner with name"BulkLoad"takeExclusiveReadLockOnRange()takes an exclusive read lock on the entire job range duringsubmitBulkLoadJob()This prevents any concurrent transactions from modifying data in the target range
Lock-aware transactions can still read from the range during the load process
The lock is automatically released via
releaseExclusiveReadLockOnRange()when the job completes, is cancelled, or errorsRange locks are managed through the
\\xff/rangeLock/keyspace with owner information in\\xff/rangeLockOwner/BulkLoad jobs will fail with
range_lock_rejectif the target range is already locked by another operation
Invariants
At any time, FDB cluster accepts at most one bulkload job.
submitBulkLoadJob()checks for existing BulkLoad or BulkDump jobs and rejects withbulkload_task_failed()if conflicts existDD partitions jobs into tasks where each task contains up to
MANIFEST_COUNT_MAX_PER_BULKLOAD_TASKmanifest entriesTask ranges are determined by the union of all manifest ranges within the job range
Tasks are assigned to data movement operations that target the appropriate storage servers for each shard
BulkLoad tasks are tracked in
BulkLoadTaskCollectionto coordinate with data movement and prevent shard boundary changesEach task validates source data integrity through manifest checksums and range validation
Tasks complete atomically - either all manifests in a task succeed or the entire task is marked as error
The job range remains exclusively locked throughout the entire operation until completion or cancellation
Task metadata persists through DD restarts - incomplete tasks are automatically resumed
Data Validation
Manifest Validation: Task ranges are validated against source manifest files using
getBulkLoadManifestMetadataFromEntry()Job Coverage Validation: The job range must be entirely covered by the source dataset or the job fails with
bulkload_dataset_not_cover_required_range()Task Atomicity: Each task either completes entirely or fails - partial task completion is not supported
SST File Integrity: Storage engines validate SST file integrity during ingestion
Range Alignment: Task ranges are aligned with shard boundaries and manifest boundaries
Failure Handling
DD Restart: Tasks persist through DD restarts via
\\xff/bulkLoadTask/metadata and are automatically resumedTask Retry: Failed tasks are retried automatically by the BulkLoad engine up to configured limits
Job Cancellation:
cancelBulkLoadJob()clears all metadata and releases range locks immediatelyData Movement Conflicts: Tasks coordinate with data movement system through
BulkLoadTaskCollectionto handle shard reassignmentsLock Conflicts: Jobs fail immediately with
range_lock_rejectif the target range is already lockedManifest Download Failures: Network/S3 failures during manifest download cause the job to error and move to history
Task Error Handling: Individual task failures are marked as
BulkLoadPhase::Errorand can be acknowledged by usersRange Coverage Failures: Jobs fail with
bulkload_dataset_not_cover_required_range()if source data doesn’t cover the requested range
Performance Considerations
Parallelism: Controlled by
DD_BULKLOAD_PARALLELISMknob for DD-level parallelismStorage Server Load: Each storage server handles one bulkload task at a time
Network Bandwidth: Large SST files may saturate network bandwidth during downloads
Storage Engine Impact: Direct SST ingestion bypasses normal write paths for better performance
Memory Usage: SST files are typically loaded into memory for validation before application
DD Pipeline Headroom:
DD_MAX_PIPELINE_MOVESmust stay well above the cluster’s peak relocation demand, or bulkload data moves are starved. See below.
DD Pipeline Headroom
DD_MAX_PIPELINE_MOVES caps the total relocations DD tracks at once, queued plus in-flight.
Bulkload is the workload most exposed to that cap, for three reasons that do not apply to ordinary data movement.
Bulkload moves are not throttled by source servers.
They read from blob storage rather than from source storage servers, so canLaunchSrc() admits them unconditionally and they consume no source work factor.
Every other relocation is self-limiting through the per-server busyness ledger; bulkload is not, which can leave a pipeline cap as the only thing gating it.
Bulkload moves run at a low priority.
They use PRIORITY_TEAM_HEALTHY (140), while shard splits use PRIORITY_SPLIT_SHARD (950).
Unlike other low-priority movement such as background rebalance, a bulkload job has a caller waiting on it.
A bulkload restore generates its own higher-priority competition. Loading a large dataset into a previously empty, contiguous range causes DD to split that range, producing thousands of split relocations that then occupy the pipeline ahead of the bulkload moves that created them.
pipelineGateActor is priority-blind, so once the cap binds, admission becomes first-come-first-served in arrival order and a bulkload move cannot be reordered ahead of the splits it caused.
This is an admission delay rather than a capacity limit: raising the cap adds no concurrency, which fetch, team and per-server budgets continue to bound, it only stops work being parked where the scheduler cannot see it.
Measured on a 1B-key restore with the cap at 1000, no bulkload data move launched for 1h54m, 43% of the run, with InQueue 893 plus InFlight 107 pinned at exactly the cap.
With headroom the first move launched in 7.15s, and in-flight moves were essentially unchanged at 107 to 114, showing the capacity to run them had been available throughout.
A 3B-key restore reached 6041 to 6815 queued relocations during the same phase, roughly six times the default of 1000.
Raise the knob before starting a large bulkload job. The default of 1000 corresponds to roughly 500 storage servers at two concurrent moves each, which a large load exceeds. Size it above the job’s expected relocation demand but below what the data distributor can hold in memory: DD resident memory grows at roughly 4 GB plus 0.11 GB per 1000 tracked relocations, and in-flight moves converge on the cap while data movement is storming, so an excessive value stops functioning as a backstop. At 100000 the distributor needs roughly 15 GB and can exhaust an 8 GB or 12 GB memory limit before the cap ever engages, whereas 20000 costs roughly 6.2 GB and leaves ample headroom over the demand measured above.
The symptom of an undersized cap is DDPipelineFull trace events during a load, with InQueue plus InFlight sitting at exactly the cap.