fix: replace the dataset once per job in overwrite-mode LanceSink - #84
Open
Zhuoxi2000 wants to merge 1 commit into
Open
Zhuoxi2000 wants to merge 1 commit into
Zhuoxi2000 wants to merge 1 commit into
Conversation
With write.mode=overwrite every LanceSink subtask deleted the dataset directory in open() and committed an Overwrite on its own first flush. With sink parallelism > 1 a subtask's first write therefore replaced the rows a peer subtask of the same job had already committed, and a subtask that opened late (or was restarted after a failover) deleted them from disk. A parallel INSERT INTO of 100 rows kept only 50. Overwrite is now decided once per job: - open() no longer deletes the dataset directory. - Commits made in overwrite mode carry the Flink job id as a Lance transaction property (flink.job-id). - On its first write a subtask replaces the dataset only if the latest version was not committed by this job; otherwise it appends. - The replacing Overwrite is committed against the version it read, so two subtasks that race to replace are detected by Lance's conflict check; the loser re-checks and appends to the winner's result. - finish() replaces the dataset with an empty one when a bounded overwrite job produced no rows, as deleting it in open() did before.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Fixes #83
What
With
write.mode=overwriteand sink parallelism > 1,LanceSinklost rows that a peer subtask had already committed. A parallelINSERT INTOof 100 rows kept 50. Both the delete inopen()and theOverwriteonisFirstWritewere decided per subtask, not per job. 88b8c42 fixed the same first-write clobber for append mode only.How
The overwrite is now decided once per job:
open()no longer deletes the dataset directory.flink.job-id.OverwritewithreadVersionset to the version it read.Overwriteas preempted ("This Overwrite transaction was preempted by concurrent transaction Overwrite"). The loser re-checks the tag and appends to the winner's result.JobIDnote under limitations).finish()handles a bounded overwrite job that produced no rows. It replaces the dataset with an empty one. Before this change, the delete inopen()removed the dataset in that case.Behaviour changes worth noting:
getRuntimeContext().getJobId(). That method is available on Flink 1.18, 1.19 and 1.20.Known limitations:
getJobId()is deprecated on 1.19 and 1.20 in favour ofgetJobInfo(), which 1.18 does not have. The sources are shared across the three modules, so I used the deprecated accessor.JobID. Flink fixes theJobIDto all zeros in application mode with HA enabled (and whenever$internal.pipeline.job-idis set), so two consecutive runs of the same application there are indistinguishable from a failover of one run, and the second run appends instead of replacing. Main always replaces in that setup, so this is a behaviour change for that deployment mode (I can add a sentence to the docs row if you want it documented). Any per-job scheme inside the legacyRichSinkFunctionhas this property, because subtasks have no coordination point other than the job id; the committer in Migrate sink to SinkV2 #49 would not.!datasetExistscreate path is unchanged from main. That is the case where the directory does not exist at all and two subtasks race to create it, which is the same window that 88b8c42 narrowed for append mode.The docs row for
overwriteininsert-into.mdnow says the replacement happens once per job.I kept this inside the legacy
LanceSink. #49 (SinkV2) would solve the problem structurally with a committer. This change does not touch the parts #49 replaces, and it merges cleanly with #82.If you would rather not add commit-protocol logic to the legacy sink before #49 lands, I am happy to cut this PR down to the minimal safe change: keep the
open()delete removal and fail fast inopen()whenwrite.mode=overwriteruns with more than one subtask, with the ITCase reduced to the parallelism check. That stops the silent loss today without the job-id tagging. Let me know which you prefer.Tests
New
LanceSinkOverwriteConcurrencyITCase, which uses the same in-JVM multi-sink pattern asLanceSinkConcurrencyITCase. Each subtask gets aMockStreamingRuntimeContextthat shares oneJobID:twoSubtasksConcurrentFirstWrite: two subtasks open, then each first-writes one row. Expects 2 rows.lateOpeningSubtaskMustNotDeletePeerRows: a subtask opens after its peer committed. Expects 2 rows.overwriteReplacesPreviousJobRows: checks that overwrite still overwrites. A seeded 3-row dataset is replaced by the job's 3 rows, and a re-run with a new job id replaces those with 1 row.overwriteWithNoRowsEmptiesDataset: a bounded job with no input leaves an empty dataset.overwriteIntoDirectoryWithoutDataset: the path exists but has no committed version yet.parallelSqlOverwriteKeepsAllRows: SQLINSERT INTOfrom datagen at parallelism 2. Expects 100 rows on two consecutive runs. The table sets ahadoop.*option only to work around thevalidateExceptempty-prefix issue that fix: adopt post-merge dataset handle, run ITCase in CI, and close type-mapping gaps #82 fixes.I also checked the race fallback separately: a stale
Overwritefalls back to append and keeps both subtasks' rows. That check is not part of this PR because it needs reflection.Before (main @ fbf4371, test file only):
After:
AI assistance: this change was drafted with an AI coding assistant (Claude) and verified locally with the tests above.