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:
- 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.
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.
- 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.
- 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:
Open questions
- Metadata column naming —
rowaddr / fragid or underscore-prefixed _rowaddr / _fragid to stay consistent with lance-spark? Flink convention leans away from leading underscores.
_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.
- Streaming commit granularity (Stage 2) — per-checkpoint commit vs. time-based batching? How to bound fragment growth?
- 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.
- 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?
- 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
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 byFragment.addColumns().ALTER TABLE a UPDATE COLUMNS c FROM view— column-level batch update without full-row rewrite, powered byFragment.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. Todaylance-flinkdoes not expose it:LanceSinkonly supportsAPPEND/OVERWRITEat whole-row granularity, and bothLanceCatalog/LanceNamespaceCatalogthrowCatalogExceptioninalterTable.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 viewto Flink SQL is tempting but a poor fit for at least four reasons: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); interactiveALTER TABLEis rare in production.ALTER TABLEin 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 executeFragment.addColumns()and write data cuts against Flink's architecture.SELECT a._rowaddr, a._fragid, b.tag_0 FROM a INNER JOIN b ON a.biz_key = b.biz_keypresumesais a bounded, static dataset. In Flink, tables default to unbounded streams; the JOIN semantics (regular/temporal/interval), and even the stability of_rowaddrunder streaming reads, all become open questions.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:
With Lance's
Fragment.updateColumns()being idempotent on_rowaddr, Flink can do something neither Spark nor MERGE INTO can:This is a genuine differentiator, and it maps naturally onto:
updateColumnson stable row addresses).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
APPEND/OVERWRITEonRowData, no column-level API_rowaddr/_fragidexposureLanceDynamicTableSourceimplementsSupportsProjectionPushDown / SupportsFilterPushDown / SupportsLimitPushDown / SupportsAggregatePushDownbut notSupportsReadingMetadataLanceCatalog.alterTableandLanceNamespaceCatalog.alterTableexplicitly throwCatalogException("...does not support altering...")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/_fragidas metadata columns (foundation)Make
LanceDynamicTableSourceimplementSupportsReadingMetadata, exposing:rowaddrBIGINT NOT NULLfragidINT NOT NULLUsage:
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 viaFragment.updateColumns().Contract:
RowDatamust carryrowaddr+fragid+ one or more payload columns._rowaddr; replay after failure is safe).Sketch:
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
DynamicTableSinkso it can be targeted by standard Flink Table APIINSERT INTO: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:
DynamicTableSinkadapter + factory option wiring on top of Stage 2's runtime.Stage 4 — Relax
LanceCatalog.alterTablefor the ADD-COLUMN case (lowest priority)Allow
ALTER TABLE t ADD COLUMN c <type>to succeed by callingDataset.addColumns()with a NULL-filled column definition (schema-only add). Keep drop / rename / type change rejected.Why last: Flink users rarely run interactive
ALTER TABLEin 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 FROMAdding this to Flink SQL requires either a Calcite parser plugin or a
SqlDialectextension, 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:
SupportsReadingMetadataforLanceDynamicTableSourceLanceColumnUpdateSink(the differentiator)DynamicTableSinkwrapper for batch backfill on FlinkLanceCatalog.alterTable— ADD COLUMN onlyOpen questions
rowaddr/fragidor underscore-prefixed_rowaddr/_fragidto stay consistent with lance-spark? Flink convention leans away from leading underscores._rowaddrstability under streaming reads — is_rowaddrguaranteed 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.lance-flinkstreaming writes.References
ALTER TABLE ... ADD COLUMNS FROMimplementationsrc/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