Skip to content

Migrate sink to SinkV2 #49

Description

@empathy87

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

Activity

  1. fightBoxing commented on Jul 23, 2026

    @fightBoxing
    Collaborator

    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 implementation Catalog / namespace implementation Closed, not merged
    #20 – Bump lance version to 5.0.0-beta.6 and arrow to 15.0.2 New 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:

    The remaining gap = this issue

    The write path is still on the legacy API:

    • LanceSink extends RichSinkFunction<RowData> implements CheckpointedFunction
    • LanceDynamicTableSink exposes it via SinkFunctionProvider.of(...)
    • Commit happens synchronously inside snapshotState() – no Committer, no 2PC, no GlobalCommitter.

    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:

    1. Keeps the append-only scope, as this issue requests.
    2. Introduces:
      • LanceSinkV2 implements TwoPhaseCommittingSink<RowData, LanceCommittable>
      • LanceCommittable (serializable: List<FragmentMetadata> + schema + writeMode) with a SimpleVersionedSerializer
      • LanceCommitter implements Committer<LanceCommittable> – the only place that calls CommitBuilder.execute(txn)
      • (optional, stretch) LanceGlobalCommitter for single-writer commit under high parallelism
    3. Wires LanceDynamicTableSink to SinkV2Provider.of(new LanceSinkV2(...)).
    4. Leaves the legacy LanceSink in place behind a write.sink-version = v1|v2 option (default v2) for one release, then removes it.
    5. 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:

    1. 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?
    2. Any preference on TwoPhaseCommittingSink vs. the newer Sink + SupportsCommitter split? I'll default to TwoPhaseCommittingSink for the smallest diff unless you object.
    3. OK to gate the switch behind write.sink-version for 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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions