All posts

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.

11 min read Hudi

TL;DR

  • Two concurrent writers against a Hudi table with the default SINGLE_WRITER both 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_control with InProcessLockProvider changed nothing: both writers committed, one writer’s data vanished, identical outcome to having no concurrency control at all.
  • That is not a bug. InProcessLockProvider is a JVM-local mutex. Two spark-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 one ConcurrentModificationException. 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_CONTROL and NON_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

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.

Buy me a coffee