Skip to content

fix: replace the dataset once per job in overwrite-mode LanceSink - #84

Open
Zhuoxi2000 wants to merge 1 commit into
lance-format:mainfrom
Zhuoxi2000:fix-overwrite-multi-subtask-clobber
Open

Zhuoxi2000 wants to merge 1 commit into
lance-format:mainfrom
Zhuoxi2000:fix-overwrite-multi-subtask-clobber

Conversation

@Zhuoxi2000

Copy link
Copy Markdown

Fixes #83

What

With write.mode=overwrite and sink parallelism > 1, LanceSink lost rows that a peer subtask had already committed. A parallel INSERT INTO of 100 rows kept 50. Both the delete in open() and the Overwrite on isFirstWrite were 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.
  • In overwrite mode, every commit carries the Flink job id as a Lance transaction property, flink.job-id.
  • On its first write, a subtask reads the latest version's transaction:
    • If this job committed it, a peer has already replaced the dataset, so the subtask appends.
    • Otherwise the subtask commits an Overwrite with readVersion set to the version it read.
  • If two subtasks race to replace the dataset, Lance rejects the second Overwrite as preempted ("This Overwrite transaction was preempted by concurrent transaction Overwrite"). The loser re-checks the tag and appends to the winner's result.
  • A failover restart keeps the job id, so a restarted subtask appends and does not wipe its peers' rows (rows replayed since the last checkpoint are appended again, which is the legacy sink's existing at-least-once behaviour; on main the restart deleted the dataset instead). A new submission of the job gets a new id and replaces the previous result (see the JobID note 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 in open() removed the dataset in that case.

Behaviour changes worth noting:

  • The previous dataset is replaced by a new version instead of being deleted from disk, so its older versions remain available for time travel until cleanup.
  • Overwrite mode now needs a runtime context, read through 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 of getJobInfo(), which 1.18 does not have. The sources are shared across the three modules, so I used the deprecated accessor.
  • The check only reads the tag on the latest version. If an unrelated writer commits to the same dataset while an overwrite job is running, a subtask whose first write lands after that commit will replace the dataset again. Main loses data in that case too, so this is not a regression.
  • The protocol is keyed on the Flink JobID. Flink fixes the JobID to all zeros in application mode with HA enabled (and whenever $internal.pipeline.job-id is 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 legacy RichSinkFunction has this property, because subtasks have no coordination point other than the job id; the committer in Migrate sink to SinkV2 #49 would not.
  • The !datasetExists create 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 overwrite in insert-into.md now 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 in open() when write.mode=overwrite runs 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 as LanceSinkConcurrencyITCase. Each subtask gets a MockStreamingRuntimeContext that shares one JobID:

  • 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: SQL INSERT INTO from datagen at parallelism 2. Expects 100 rows on two consecutive runs. The table sets a hadoop.* option only to work around the validateExcept empty-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 Overwrite falls 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):

JAVA_HOME=/opt/homebrew/opt/openjdk@17 mvn -B -ntp -am -pl lance-flink-1.20 test \
  -Dtest=LanceSinkOverwriteConcurrencyITCase -Dsurefire.failIfNoSpecifiedTests=false

Tests run: 6, Failures: 4, Errors: 1
  twoSubtasksConcurrentFirstWrite         expected: 2L   but was: 1L
  lateOpeningSubtaskMustNotDeletePeerRows expected: 2L   but was: 1L
  overwriteReplacesPreviousJobRows        expected: 3L   but was: 2L
  parallelSqlOverwriteKeepsAllRows        expected: 100L but was: 50L
  overwriteWithNoRowsEmptiesDataset       IllegalArgument: Dataset at path ... was not found

After:

# lance-flink-1.20
-Dtest=LanceSinkOverwriteConcurrencyITCase,LanceSinkConcurrencyITCase,LanceSinkTest   -> Tests run: 18, Failures: 0, Errors: 0
-Dtest=LanceConnectorITCase,LanceSqlITCase,LanceUpsertSinkITCase                     -> 15 + 20 + 8 run, 0 failures
# lance-flink-1.18 and lance-flink-1.19
-Dtest=LanceSinkOverwriteConcurrencyITCase,LanceSinkConcurrencyITCase                -> 7 run, 0 failures (each module)

AI assistance: this change was drafted with an AI coding assistant (Claude) and verified locally with the tests above.

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.
@github-actions github-actions Bot added the bug Something isn't working label Oct 2, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

write.mode=overwrite LanceSink with parallelism > 1 drops rows committed by peer subtasks

1 participant