Repository navigation
Migrate sink to SinkV2 #49
Description
Activity
Happy to pick this up. Sharing a short status check on the two prerequisites and a concrete plan for the SinkV2 migration.
Prerequisite status (as of today)
PR Actual scope State #14 – feat: add catalog namespace support and refactor adapter implementationCatalog / namespace implementation Closed, not merged #20 – Bump lance version to 5.0.0-beta.6 and arrow to 15.0.2New Lance API + Arrow migration Open, not merged Minor: the issue description swapped the two links – #14 is the catalog one, #20 is the API/Arrow bump. Not a blocker, just for future readers.
So on this upstream repo, neither prerequisite has actually landed yet.
Where we are downstream
We have been carrying both prerequisites out-of-tree and they now work end-to-end:
- Catalog / namespace – implemented on top of
org.lance.namespace.LanceNamespace(CreateNamespaceRequest,CreateTableRequest,ListTables…, etc.) via aLanceNamespaceCatalog+LanceNamespaceCatalogFactorypair. Functionally equivalent to what feat: add catalog namespace support and refactor adapter implementation #14 was trying to introduce. - New Lance API + Arrow – all read/write paths already use
CommitBuilder/Transaction/Append/Overwrite/FragmentMetadata/WriteParamsfrom the newerorg.lance.*surface (equivalent to the target state of Bump lance version to 5.0.0-beta.6 and arrow to 15.0.2 #20).
The remaining gap = this issue
The write path is still on the legacy API:
LanceSink extends RichSinkFunction<RowData> implements CheckpointedFunctionLanceDynamicTableSinkexposes it viaSinkFunctionProvider.of(...)- Commit happens synchronously inside
snapshotState()– noCommitter, no 2PC, noGlobalCommitter.
This is exactly the sync-phase-commit shape that has been reported to cause duplicate rows under concurrent compaction on S3, so migrating to SinkV2 is worth doing on its own merits, independent of #14/#20.
Proposal – I'd like to own this
I can send a PR that:
- Keeps the append-only scope, as this issue requests.
- Introduces:
LanceSinkV2 implements TwoPhaseCommittingSink<RowData, LanceCommittable>LanceCommittable(serializable:List<FragmentMetadata> + schema + writeMode) with aSimpleVersionedSerializerLanceCommitter implements Committer<LanceCommittable>– the only place that callsCommitBuilder.execute(txn)- (optional, stretch)
LanceGlobalCommitterfor single-writer commit under high parallelism
- Wires
LanceDynamicTableSinktoSinkV2Provider.of(new LanceSinkV2(...)). - Leaves the legacy
LanceSinkin place behind awrite.sink-version = v1|v2option (defaultv2) for one release, then removes it. - Adds an IT that runs a checkpointed streaming append against a local dataset and asserts exactly-once (no duplicates on TM restart).
Questions for maintainers before I open the PR:
- Are feat: add catalog namespace support and refactor adapter implementation #14 and Bump lance version to 5.0.0-beta.6 and arrow to 15.0.2 #20 still the intended path, or has catalog / API-bump work been superseded by another PR I should rebase on? If feat: add catalog namespace support and refactor adapter implementation #14 is dead, would you accept a fresh catalog PR alongside the SinkV2 one, or should SinkV2 land first?
- Any preference on
TwoPhaseCommittingSinkvs. the newerSink+SupportsCommittersplit? I'll default toTwoPhaseCommittingSinkfor the smallest diff unless you object. - OK to gate the switch behind
write.sink-versionfor one release, or would you rather flip the default immediately with no v1 fallback?
Happy to self-assign if that's OK, and to pair on design.
- Catalog / namespace – implemented on top of
I would like to migrate writes in the Lance Flink connector to the Flink SinkV2 API.
This issue is only about append mode.
Before this can be done, it would be better to have the catalog implementation in place and migration to new Lance API and Arrow:
#20
#14