Skip to content

Bring Lance's low-cost column evolution to the Flink connector (streaming-first, staged proposal) #61

Description

@fightBoxing

Background

The lance-spark connector recently exposed two SQL extensions built on top of Lance's native data-evolution APIs:

  • ALTER TABLE a ADD COLUMNS c FROM view — add & backfill a new column without rewriting existing rows, powered by Fragment.addColumns().
  • ALTER TABLE a UPDATE COLUMNS c FROM view — column-level batch update without full-row rewrite, powered by Fragment.updateColumns().

Reference: Lance + Spark SQL: 实现低成本加列回填数据与修改列 (official Lance & LanceDB post).

The underlying capability — column-file independence + row-address (_rowaddr + _fragid) addressing — lives in the Lance format layer, not in Spark. It should be reachable from Flink too. Today lance-flink does not expose it: LanceSink only supports APPEND / OVERWRITE at whole-row granularity, and both LanceCatalog / LanceNamespaceCatalog throw CatalogException in alterTable.

This issue proposes bringing the same underlying data-evolution capability to the Flink connector — but adapted to Flink's stream-first semantics, not by porting Spark's SQL surface.

Why not port lance-spark's SQL syntax to Flink

Before proposing anything, we want to be explicit about what we are not doing and why. A naive port of ALTER TABLE ... ADD/UPDATE COLUMNS FROM view to Flink SQL is tempting but a poor fit for at least four reasons:

  1. User habits differ. Spark users routinely run interactive DDL in spark-sql / Zeppelin / notebooks — one-shot jobs, inspect result, next statement. Flink SQL is primarily used to declare long-running topologies (CREATE TABLE + INSERT INTO); interactive ALTER TABLE is rare in production.
  2. ALTER TABLE in Flink is a metadata operation, not a data operation. It updates the Catalog but does not trigger a data-mutation job. Wiring a DDL statement to actually execute Fragment.addColumns() and write data cuts against Flink's architecture.
  3. The batch pattern in the article assumes a finite target table. SELECT a._rowaddr, a._fragid, b.tag_0 FROM a INNER JOIN b ON a.biz_key = b.biz_key presumes a is a bounded, static dataset. In Flink, tables default to unbounded streams; the JOIN semantics (regular/temporal/interval), and even the stability of _rowaddr under streaming reads, all become open questions.
  4. If a user just wants one-shot batch backfill, they should use Spark. Duplicating Spark's exact interactive experience in Flink would give users no reason to choose Flink.

So the design goal here is not "match lance-spark's SQL". It is: expose Lance's low-cost column evolution as a first-class primitive in the Flink runtime, then build the pipelines Flink can uniquely deliver on top.

Why this matters (Flink-specific value)

The one thing Flink can do that Spark cannot: streaming column backfill.

Consider an ML feature pipeline where new feature values arrive continuously (Kafka, CDC, upstream micro-service). Today's options are all bad:

  • Rewrite the whole target table on every batch → prohibitive I/O.
  • MERGE INTO on every micro-batch → whole-row rewrite, high write amplification.
  • Buffer offline and run Spark daily → loses freshness.

With Lance's Fragment.updateColumns() being idempotent on _rowaddr, Flink can do something neither Spark nor MERGE INTO can:

Kafka(feature updates) → Flink stream → JOIN Lance dim view to get _rowaddr
                                     → per-checkpoint column-only commit → Lance

This is a genuine differentiator, and it maps naturally onto:

  • Flink's exactly-once semantics (idempotent updateColumns on stable row addresses).
  • The existing time-travel and vector-search story already shipped in this repo.
  • ML feature freshness requirements that offline Spark cannot meet.

Batch feature backfill inside Flink is a secondary, opportunistic use case — worth supporting for users already on Flink, but not the primary justification.

Current gaps in lance-flink

Area Status
Sink write modes Only APPEND / OVERWRITE on RowData, no column-level API
_rowaddr / _fragid exposure Not exposed — LanceDynamicTableSource implements SupportsProjectionPushDown / SupportsFilterPushDown / SupportsLimitPushDown / SupportsAggregatePushDown but not SupportsReadingMetadata
Catalog Both LanceCatalog.alterTable and LanceNamespaceCatalog.alterTable explicitly throw CatalogException("...does not support altering...")
Lance Java SDK bindings Fragment.addColumns() / Fragment.updateColumns() are available upstream but never referenced in this repo (verified: grep -r "addColumns|updateColumns" → 0 hits)

Proposed design (staged)

Four stages, ordered by Flink-fit rather than by mechanical parity with Spark. Each stage is independently shippable.

Stage 1 — Expose _rowaddr / _fragid as metadata columns (foundation)

Make LanceDynamicTableSource implement SupportsReadingMetadata, exposing:

Metadata key Flink type Semantics
rowaddr BIGINT NOT NULL Lance physical row address
fragid INT NOT NULL Fragment ID

Usage:

CREATE TABLE lance_users (
    id BIGINT,
    name STRING,
    rowaddr BIGINT METADATA FROM 'rowaddr' VIRTUAL,
    fragid  INT    METADATA FROM 'fragid'  VIRTUAL
) WITH ('connector' = 'lance', 'path' = '...');

Why first: pure additive change to a single source class; unblocks every downstream stage; also useful on its own for row-level upsert and diagnostic queries.
Scope: single class change + IT case. Low risk.

Stage 2 — Streaming LanceColumnUpdateSink (Flink's main differentiator)

Checkpoint-integrated sink that consumes a stream of (rowaddr, fragid, newValue…) records and commits column-only updates per checkpoint via Fragment.updateColumns().

Contract:

  • Input RowData must carry rowaddr + fragid + one or more payload columns.
  • Idempotent under exactly-once (Lance keys updates on _rowaddr; replay after failure is safe).
  • Commits are batched at checkpoint boundaries to bound fragment growth.

Sketch:

DataStream<RowData> featureUpdates = /* Kafka → JOIN lance dim view → project (rowaddr, fragid, v) */;
featureUpdates.sinkTo(LanceColumnUpdateSink.forTable("lance_users")
        .targetColumns("tag_0")
        .mode(UPDATE));

Why this is the headline stage: this pipeline cannot be expressed by Spark. It is the concrete, differentiated value of bringing Lance data evolution to Flink.
Scope: new sink + checkpoint/commit protocol + IT case with fault-injection. Highest engineering investment, highest strategic value.

Stage 3 — Batch LanceColumnUpdateSink (opportunistic, for users already on Flink)

Same sink class as Stage 2 but exposed as a DynamicTableSink so it can be targeted by standard Flink Table API INSERT INTO:

Table view = tEnv.sqlQuery(
    "SELECT a.rowaddr, a.fragid, b.tag_0 " +
    "FROM lance_users a INNER JOIN feature_source b ON a.id = b.biz_key");

view.executeInsert("lance_column_sink");   // configured with mode=ADD, target-columns=tag_0

Positioning: this covers the same use case as lance-spark's ADD/UPDATE COLUMNS FROM, but for users whose whole pipeline is already on Flink and don't want to spin up Spark just to backfill one column. Not a reason to choose Flink over Spark on its own; a nice-to-have.
Scope: DynamicTableSink adapter + factory option wiring on top of Stage 2's runtime.

Stage 4 — Relax LanceCatalog.alterTable for the ADD-COLUMN case (lowest priority)

Allow ALTER TABLE t ADD COLUMN c <type> to succeed by calling Dataset.addColumns() with a NULL-filled column definition (schema-only add). Keep drop / rename / type change rejected.

Why last: Flink users rarely run interactive ALTER TABLE in production; the value is mostly cosmetic parity with other catalogs.
Scope: 2 catalog classes, gated by capability flag.

(Explicitly not proposed) SQL syntax extension ALTER TABLE ... ADD/UPDATE COLUMNS FROM

Adding this to Flink SQL requires either a Calcite parser plugin or a SqlDialect extension, plus multi-version maintenance across Flink 1.18/1.19/1.20. Given (a) the arguments above about engine fit and (b) Stage 3 already covering the same use case via Table API, we do not plan this and would push back on it unless there's clear community demand.

Scope of this issue

Umbrella / discussion ticket. If direction is accepted, each stage becomes its own tracked issue:

  • Stage 1 — SupportsReadingMetadata for LanceDynamicTableSource
  • Stage 2 — Streaming LanceColumnUpdateSink (the differentiator)
  • Stage 3 — DynamicTableSink wrapper for batch backfill on Flink
  • Stage 4 — LanceCatalog.alterTable — ADD COLUMN only
  • (Deferred) SQL syntax extension

Open questions

  1. Metadata column naming — rowaddr / fragid or underscore-prefixed _rowaddr / _fragid to stay consistent with lance-spark? Flink convention leans away from leading underscores.
  2. _rowaddr stability under streaming reads — is _rowaddr guaranteed stable across snapshots between the JOIN read and the column-update commit? This is the correctness cornerstone of Stage 2 and needs an authoritative answer from the Lance team.
  3. Streaming commit granularity (Stage 2) — per-checkpoint commit vs. time-based batching? How to bound fragment growth?
  4. Compaction interaction (Stage 2) — do we auto-schedule compaction inside the sink, or expose a knob and let users run compaction externally? This is a known pre-existing pain point for lance-flink streaming writes.
  5. Nullability semantics (Stage 3) — for ADD-column-from-view when some target rows have no matching source row, match Spark's INNER-JOIN-defaults-to-NULL behavior?
  6. Type validation — plan-time schema check on target column vs. payload column, or runtime error?

References

  • Article: Lance + Spark SQL: 实现低成本加列回填数据与修改列
  • lance-spark: ALTER TABLE ... ADD COLUMNS FROM implementation
  • Lance format docs — data evolution & row addresses
  • Related in-repo files:
    • src/main/java/org/apache/flink/connector/lance/table/LanceDynamicTableSource.java (Stage 1 landing site)
    • src/main/java/org/apache/flink/connector/lance/LanceSink.java (Stage 2/3 template)
    • src/main/java/org/apache/flink/connector/lance/table/LanceCatalog.java (Stage 4 landing site)

Happy to send Stage 1 as a first small PR to validate the approach.

cc @LuciferYang @jackye1995

Activity

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

    enhancementNew feature or request

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions