Concurrency control in Hudi: the lock provider that protects nothing
Two writers, one Hudi table, three configurations, measured. The default loses an entire commit silently. Turning on optimistic concurrency control with the lock provider that needs no infrastructure loses it just as silently. Only one of the three tells the losing writer that it lost.
- What “concurrent” means here
- The default loses a commit, silently
- Optimistic concurrency, and the provider that does nothing
- A lock that actually spans writers
- The three runs side by side
- What OCC actually does
- Choosing a lock provider
- Table services are writers too
- Non-blocking concurrency control
- Common misconceptions
- Frequently asked questions
- Conclusion
- References
- Trademarks
TL;DR
- Two concurrent writers against a Hudi table with the default
SINGLE_WRITERboth reported success, both commits landed on the timeline, and the table contained only one writer’s rows. The other 200 updates were gone, with no error anywhere.- Turning on
optimistic_concurrency_controlwithInProcessLockProviderchanged nothing: both writers committed, one writer’s data vanished, identical outcome to having no concurrency control at all.- That is not a bug.
InProcessLockProvideris a JVM-local mutex. Twospark-submits are two JVMs, so neither ever sees the other’s lock — and it is the provider you reach for precisely because it needs no infrastructure.- With
ZookeeperBasedLockProvider, the same test produced one commit and oneConcurrentModificationException. The losing writer knows it lost, which is the entire point.- OCC is not a lock held for the duration of a write. Writers work freely and take a short lock at commit time to check whether what they touched moved underneath them.
- Hudi 1.1.1 has three modes —
SINGLE_WRITER,OPTIMISTIC_CONCURRENCY_CONTROLandNON_BLOCKING_CONCURRENCY_CONTROL.
The Merge-on-Read post notes in passing that two concurrent jobs “can both succeed and one can lose its changes”. That sentence is easy to read past. This post runs it.
Three configurations, two genuinely concurrent spark-submit processes each time,
both writing the same 200 record keys into the same table. The only thing that
changes between runs is the concurrency configuration.
What “concurrent” means here
Hudi coordinates writers through its timeline. A write creates an instant, moves
it through requested and inflight, and completes it. Two writers overlapping
in time are only a problem when they touch the same file groups — and with
upserts on the same keys, they always do.
The test is therefore the worst realistic case: same table, same keys, same partition, overlapping in time. Each writer is a separate JVM, because that is what two Spark applications are.
df = [(i, f"w{writer_id}", writer_id) for i in range(1, 201)]
df.write.format("hudi").options(**opts).mode("append").save(PATH)
The default loses a commit, silently
hoodie.write.concurrency.mode defaults to SINGLE_WRITER. The name is a
statement about what Hudi assumes, not a restriction it enforces — nothing stops
a second writer, and nothing tells you one arrived.
RESULT writer=1 mode=single OUTCOME=committed
RESULT writer=2 mode=single OUTCOME=committed
Both succeeded. Both are on the timeline, four milliseconds apart:
20260928043109150_20260928043355742.commit
20260928043109154_20260928043355529.commit
And the table:
TOTAL_ROWS 200
BY_WRITER [('w2', 200)]
COMMITS ['109154']
Writer 1’s commit exists on the timeline and none of its data is visible. Two hundred updates were acknowledged to the application that made them and then discarded. No exception, no warning, no failed job to alert on — the only way to discover this is to notice the rows are wrong.
That is the failure mode worth internalising, because every layer above Hudi reports success.
Optimistic concurrency, and the provider that does nothing
The documented fix is optimistic concurrency control:
"hoodie.write.concurrency.mode": "optimistic_concurrency_control",
"hoodie.write.lock.provider": "org.apache.hudi.client.transaction.lock.InProcessLockProvider",
"hoodie.cleaner.policy.failed.writes": "LAZY",
InProcessLockProvider is the one that appears in every quick-start, because it
is the only one that needs no infrastructure. Running the identical test:
RESULT writer=1 mode=occ_inproc OUTCOME=committed
RESULT writer=2 mode=occ_inproc OUTCOME=committed
TOTAL_ROWS 200
BY_WRITER [('w1', 200)]
Exactly the same failure. Both writers committed, one writer’s two hundred rows are gone, nothing raised an error. Enabling optimistic concurrency control bought precisely nothing.
The reason is in the class name and is not hidden: the lock is held in process.
It is a mutex inside one JVM. Two spark-submit invocations are two JVMs, so
writer 1’s lock is invisible to writer 2 and vice versa. Each acquires “the” lock
instantly, sees no conflict, and commits.
flowchart TB
subgraph J1["JVM 1"]
L1["InProcessLockProvider<br/>lock acquired"]
end
subgraph J2["JVM 2"]
L2["InProcessLockProvider<br/>lock acquired"]
end
L1 --> T["one table"]
L2 --> T
T --> X["both commit<br/>one loses"]
This is the trap the post is named for. The configuration looks correct — the mode is set, a provider is named, the job logs show OCC is on — and it protects nothing across the only boundary that matters. It is genuinely useful for concurrent writers inside a single application, and misleading everywhere else.
A lock that actually spans writers
Swapping the provider for one backed by shared infrastructure, with everything else identical:
"hoodie.write.concurrency.mode": "optimistic_concurrency_control",
"hoodie.write.lock.provider": "org.apache.hudi.client.transaction.lock.ZookeeperBasedLockProvider",
"hoodie.write.lock.zookeeper.url": "zookeeper",
"hoodie.write.lock.zookeeper.port": "2181",
"hoodie.write.lock.zookeeper.base_path": "/hudi_locks",
"hoodie.write.lock.zookeeper.lock_key": "cc_demo",
"hoodie.cleaner.policy.failed.writes": "LAZY",
RESULT writer=2 mode=occ_zk OUTCOME=committed
RESULT writer=1 mode=occ_zk OUTCOME=failed(ConcurrentModificationException)
TOTAL_ROWS 200
BY_WRITER [('w2', 200)]
One writer committed. The other failed loudly with
ConcurrentModificationException. The table holds a consistent result, and the
writer that lost knows it lost, so it can retry.
The lock was real — ZooKeeper shows the znode:
ls /hudi_locks -> [cc_demo]
Note that the table contents are the same shape as the broken runs: one writer’s 200 rows. The difference is not in the data, it is in whether anybody was told. That is the whole value of concurrency control, and it is why “the data looked fine” is not evidence that a configuration works.
The three runs side by side
| Mode | Lock provider | Writer 1 | Writer 2 | Table | Silent loss |
|---|---|---|---|---|---|
SINGLE_WRITER |
— | committed | committed | only w2 | yes |
optimistic_concurrency_control |
InProcessLockProvider |
committed | committed | only w1 | yes |
optimistic_concurrency_control |
ZookeeperBasedLockProvider |
failed | committed | only w2 | no |
Two of the three configurations lose data without telling anyone, and one of those two is the configuration people believe protects them.
What OCC actually does
“Optimistic” is precise, and worth understanding before enabling it.
Writers do not hold a lock while they work. Each writes its files freely, and only at commit time takes a short lock to check whether the file groups it touched were modified by a commit that landed while it was working. If they were, it aborts.
flowchart LR
W["write files<br/>(no lock)"] --> L["take lock"]
L --> C{"did anything<br/>I touched change?"}
C -->|no| OK["commit, release"]
C -->|yes| AB["abort<br/>ConcurrentModificationException"]
Two consequences follow:
- Cheap when collisions are rare, expensive when they are common. Two writers touching disjoint partitions rarely conflict and rarely retry. Two writers upserting the same keys — this test — conflict every time, and the work of the loser is thrown away entirely. If your writers always collide, OCC converts a correctness problem into a throughput problem; partitioning the work between writers is the better answer.
- The retry is yours. Hudi raises the exception; it does not re-run the write. A pipeline enabling OCC needs a retry path, and that retry must be safe to run twice.
hoodie.cleaner.policy.failed.writes set to LAZY belongs with OCC for a related
reason: with multiple writers, the eager policy would let one writer clean up
files belonging to another writer’s in-flight commit. Lazy cleaning defers that to
the cleaner, which can tell an abandoned write from a running one.
Choosing a lock provider
The bundle ships several, and the choice is about where the lock is visible:
| Provider | Lock lives in | Use when |
|---|---|---|
InProcessLockProvider |
one JVM’s memory | concurrent writers inside one application; tests |
FileSystemBasedLockProvider |
a file on the table’s filesystem | simple setups on a filesystem with atomic create-if-absent semantics |
StorageBasedLockProvider |
object storage | object stores offering the conditional writes this needs |
ZookeeperBasedLockProvider |
ZooKeeper | you already run ZooKeeper — as anyone running Kafka does |
HiveMetastoreBasedLockProvider |
the Hive Metastore | you already run one, which most Hudi deployments do |
DynamoDBBasedLockProvider |
DynamoDB | AWS, no ZooKeeper or metastore to lean on |
The two worth noticing are the ends of the list. InProcessLockProvider
protects less than people assume, and the metastore or ZooKeeper providers
usually need no new infrastructure, because a Hudi deployment already has one
of them. The gap between the broken configuration and the correct one is
frequently just a class name and three connection properties.
A caution on the filesystem provider: it depends on the underlying store giving atomic create-if-absent, which not every object store does with the semantics a lock needs. Prefer a provider whose backing store was designed for coordination.
Table services are writers too
Compaction, clustering and cleaning modify the table, so they participate in all of the above. An asynchronous compaction running while an ingest job writes is exactly the two-writer situation in this post.
The common arrangement that avoids most of the difficulty: one writer per table, with table services running inline or from a single dedicated job. It is less flexible and it removes the whole class of problem. Where genuinely concurrent ingestion is required, enable OCC with a real lock provider and accept the retries.
Merge-on-Read tolerates concurrent ingestion better than Copy-on-Write, because two writers appending log files to the same file group are not necessarily overwriting each other. It does not remove the need to coordinate — compaction still must be — and it is a reason MoR is the usual choice for multi-writer ingest.
Non-blocking concurrency control
Hudi 1.x adds a third mode. The enum in 1.1.1 confirms it:
SINGLE_WRITER
OPTIMISTIC_CONCURRENCY_CONTROL
NON_BLOCKING_CONCURRENCY_CONTROL
NBCC is designed to let concurrent writers proceed without taking a lock at all, resolving ordering at read and compaction time instead — which is what the completion-time ordering of the 1.x timeline exists to support, and the same mechanism behind the incremental-read behaviour covered separately. It carries requirements on table type and index that OCC does not.
To be explicit about the limits of this post: the enum above is verified; NBCC’s behaviour is not. The three measured scenarios here are all OCC and single-writer. Treat NBCC as worth investigating against your own table type and index rather than as something demonstrated.
Common misconceptions
“SINGLE_WRITER prevents a second writer.” It is an assumption Hudi makes,
not a constraint it enforces. A second writer proceeds and one of them loses.
“I enabled OCC, so I am safe.” Measured above: OCC with
InProcessLockProvider across two Spark applications lost data exactly as the
default did.
“OCC holds a lock while writing.” It holds a short lock at commit time only.
“A conflict means Hudi retries for me.” It raises
ConcurrentModificationException and stops. Retrying is your job.
“The data looked right, so the configuration works.” Two of the three runs here produced a plausible-looking table. Correctness of the contents is not evidence that concurrency control is functioning.
“Lock providers need new infrastructure.” ZooKeeper and the Hive Metastore are already present in most Hudi deployments.
Frequently asked questions
What exactly counts as a conflict? Overlapping file groups between the concurrent commits, under the default conflict-resolution strategy. Writers in different partitions usually do not conflict; writers upserting the same keys always do.
How do I know OCC is actually working?
Force a conflict and confirm one writer fails. If both succeed against the same
keys, it is not working — which is precisely how the InProcessLockProvider
configuration above looks when it is doing nothing.
Should I just avoid multiple writers? If you can, yes. One writer per table with table services in a single job removes the whole class of problem, and is the right default for most pipelines.
Does this apply to Iceberg and Delta too? The concept does; the mechanism differs. Iceberg commits by swapping a pointer in the catalog, so the catalog provides the atomicity — see Iceberg’s write path, where two concurrent writers produced a linear snapshot chain with no lost rows and no external lock.
What about the metadata table? It is itself a Hudi table and is written as part of your commits, so it is subject to the same coordination. Multi-writer setups should have concurrency control configured correctly before enabling anything that increases write contention.
Conclusion
The numbers worth carrying away are that two of three configurations lost two
hundred rows without raising anything, and that the one people actually deploy —
optimistic_concurrency_control with InProcessLockProvider — was one of them.
The underlying lesson is about how this class of bug presents. There was no exception, no failed job, no alert. Both writers were told they had committed. The only symptom was that the table contained the wrong rows, which nothing was checking. A configuration that protects nothing looks exactly like one that works, until you deliberately force a conflict and watch whether anybody complains.
So the test is the deliverable here, more than the settings: write the same keys from two separate processes and confirm that one of them fails. If both succeed, you do not have concurrency control, whatever the configuration says.
References
- Hudi concurrency control for the modes, lock providers and their configuration
- Hudi configurations for
hoodie.write.lock.*and the cleaner policies that accompany OCC WriteConcurrencyMode.javafor the three modesInProcessLockProvider.javafor what its scope actually is- Hudi Merge-on-Read for why MoR tolerates concurrent ingest better
- Incremental processing in Hudi for the completion-time timeline ordering NBCC builds on
Trademarks
Apache Hudi, Apache Spark, Apache ZooKeeper, Apache Hive, Apache Iceberg, Apache and the Apache feather logo are either registered trademarks or trademarks of The Apache Software Foundation in the United States and other countries. Amazon DynamoDB is a trademark of Amazon.com, Inc. or its affiliates.
Found this useful?
These posts and tools are free. If one saved you an afternoon, you can buy me a coffee.