[LIVY-1078] Add StateStore.tryExclusiveCreate() abstraction across state-store impls - #552
soumyadeeplogin wants to merge 3 commits into
Conversation
Adds a create-if-absent primitive to the StateStore abstraction so callers can atomically persist a value only when no value is already stored at that key, instead of always overwriting via set(). - ZooKeeperStateStore/ZooKeeperManager: relies on ZooKeeper's create() failing with NodeExistsException, avoiding a racy separate exists check. - FileSystemStateStore: uses FileContext.create() with CREATE (no OVERWRITE), which fails atomically with FileAlreadyExistsException. - BlackholeStateStore: always returns true (no-op store, nothing to conflict with). - SessionStore.trySave() exposes the primitive at the session level.
|
Filed LIVY-1078 for this change. |
|
Hi @gyogal — this one is approved and CI is green. Could you merge it when you have a chance? Thanks! |
|
Thanks for your contribution @soumyadeeplogin , it looks good, however I have one additional question before merging: the original |
tryExclusiveCreate() previously wrote directly to the destination path and relied on CreateFlag.CREATE (no OVERWRITE) to fail atomically when the key already existed. If livy-server crashed mid-write, the destination would be left partially written yet "claimed" forever, since every future exclusive-create attempt would see the existing file and back off. Write to a uniquely-named temp file first, then perform the exclusive claim via a non-overwriting atomic rename (mirrors set()'s existing write-then-rename shape). The destination is only ever visible once it is complete, and the rename itself still fails atomically with FileAlreadyExistsException if another writer wins the race, so no separate (and racy) exists-check is needed. Clean up the orphaned temp file when the rename loses the race.
|
Good catch, thanks for looking closely — you're right, and it was not intentional. The direct create() at the final path meant a crash between create and close would leave a truncated file at the destination, and every future tryExclusiveCreate call for that key would see it as already claimed (FileAlreadyExistsException) with no way to recover. Pushed a fix: tryExclusiveCreate() now writes to a uniquely-named temp file first, then performs the exclusive claim via a non-overwriting atomic rename (fileContext.rename(tmpPath, absPath(key), Rename.NONE)) — same write-then-rename shape set() already uses. The destination is only ever visible once it's complete, and the rename itself still fails atomically with FileAlreadyExistsException if another writer wins the race, so the exclusivity guarantee holds. The temp file (and its .crc sidecar) is cleaned up if the rename loses the race. Updated the FileSystemStateStoreSpec tests to cover this and re-ran the full suite locally — 14/14 pass. Let me know if you'd like anything else adjusted. |
Wrap the rename in a try/finally so the uniquely-named temp file is deleted whenever the rename doesn't succeed, not just when it fails with FileAlreadyExistsException. Unlike set()'s fixed-name temp file, tryExclusiveCreate's UUID-named temp file has no future call that would reclaim and clean it up, so any other exception during the rename previously leaked it.
|
@soumyadeeplogin Thanks for updating your PR. Could you please update the description so that the "Was this patch authored or co-authored using generative AI tooling?" section has the "Generated-by:" header, similarly to the recent Livy or Spark commits? |
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## master #552 +/- ##
============================================
+ Coverage 68.68% 68.94% +0.26%
- Complexity 1218 1239 +21
============================================
Files 106 107 +1
Lines 6815 7014 +199
Branches 836 862 +26
============================================
+ Hits 4681 4836 +155
- Misses 1666 1703 +37
- Partials 468 475 +7 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
@gyogal Updated the description to include the `Generated-by:` header, matching the convention from #546 — thanks for flagging it. Also, per your earlier comment on the temp file: pushed a `try`/`finally` around the rename in `tryExclusiveCreate()` so the temp file is cleaned up on any rename failure, not just `FileAlreadyExistsException`. Added a test for that path too, 15/15 passing. Let me know if this looks good — appreciate an approve/merge when you have a chance. |
| def tryCreate(key: String, value: Object): Boolean = { | ||
| val data = serializeToBytes(value) | ||
| try { | ||
| curatorClient.create().creatingParentsIfNeeded().forPath(key, data) |
There was a problem hiding this comment.
I asked Claude to take a closer look and a potential correctness issue was pointed out here:
Fix before this gets a caller — ZooKeeperManager.tryCreate:
Under the retry policy (RetryNTimes), a create() whose ack is lost can be retried and
throw a spurious NodeExistsException for the node this client just created — so
tryCreate returns false ("someone else won") when it actually won. set() shares the
underlying create() but doesn't treat NodeExistsException as a signal, so this is the
first place it becomes a correctness bug.
.withProtection() is not the fix here — it mangles the znode name (_c_<uuid>-…),
breaking this store's key == path contract. Disambiguate by reading the node back instead:
} catch {
- case _: NodeExistsException => false
+ case _: NodeExistsException =>
+ // Under the retry policy, a create() whose ack was lost can be retried and see the
+ // node this client itself just created. Read back to disambiguate: if it holds the
+ // bytes we wrote, we are the creator. (Sound while values are unique per creator —
+ // true for session RecoveryMetadata, which embeds the id.)
+ java.util.Arrays.equals(curatorClient.getData().forPath(key), data)
}@soumyadeeplogin Please let me know if you think this is worth considering before tryCreate() gets used. I am not familiar with this use case and I can't guarantee that the above code change is correct, it definitely needs to be tested. If this change is added, new unit tests could be useful for this scenario.
There was a problem hiding this comment.
Actually this fix may not be correct if two threads want to create the node with the same content and whether this fix is needed depends on how tryCreate() will be used. Your change may be OK in its current form if you feel like this is not an issue that can happen in a real life scenario, so please let me know if any fix is needed to the existing logic.
What changes were proposed in this pull request?
Adds
StateStore.tryExclusiveCreate(key, value): Boolean, an atomic create-if-absent primitive alongside the existingset()(which always overwrites). Each backend implements it using its own native atomic guarantee rather than a check-then-act race:create()failing withNodeExistsExceptionif the znode already exists.FileContext.create()withCREATE(noOVERWRITE), which fails atomically withFileAlreadyExistsException.true(recovery disabled, so there's nothing to conflict with).SessionStore.trySave()exposes this at the session level for callers that need "create this session's recovery record only if it doesn't already exist" semantics.This is plumbing only — the new method is unused by any existing caller in this PR, so there is no behavior change. It's a building block for a follow-up fix to a session-ID collision issue (a stale id counter on a recovered Livy instance can otherwise cause one session to silently overwrite another's state file).
How was this patch tested?
Added unit tests to
ZooKeeperStateStoreSpec,FileSystemStateStoreSpec,BlackholeStateStoreSpec, andSessionStoreSpec, covering both the create-succeeds path and the key-already-exists path (asserting no side effects / no overwrite) for every backend.Ran
mvn -pl server -am -Pscala-2.12 -Pspark3 verify: 55/55 tests pass (8 new), scalastyle/checkstyle/RAT all clean.Also addressed review feedback:
FileSystemStateStore.tryExclusiveCreate()now writes to a uniquely-named temp file and performs the exclusive claim via a non-overwriting atomic rename (crash-safety fix), and wraps that rename in atry/finallyso the temp file is cleaned up on any rename failure, not just a lost claim race.Was this patch authored or co-authored using generative AI tooling?
Yes. Generated-by: Claude Code (Sonnet 5, Anthropic), used to implement
tryExclusiveCreate()across allStateStorebackends and their corresponding unit tests, and to draft/iterate on the follow-up crash-safety and temp-file-cleanup fixes from review feedback, under human review throughout. Please refer to the ASF Generative Tooling Guidance for details.