Skip to content

[SPARK-56132][SS] Call pruneColumns on V2 streaming to fix metadata reading issue - #56133

Closed
zikangh wants to merge 15 commits into
apache:masterfrom
zikangh:stack/prunecolumns-streaming
Closed

zikangh wants to merge 15 commits into
apache:masterfrom
zikangh:stack/prunecolumns-streaming

Conversation

@zikangh

@zikangh zikangh commented May 27, 2026 •

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

In MicroBatchExecution.logicalPlan, before calling build() on the V2 streaming scan
builder, call SupportsPushDownRequiredColumns.pruneColumns(output.toStructType) if the
builder supports it. output is the analyzed relation output, which already includes any
metadata columns the query references (added by the AddMetadataColumns rule).

Why are the changes needed?

  1. Metadata column reads in V2 streaming crash with ArrayIndexOutOfBoundsException.
    When a query selects a metadata column (e.g. _metadata.row_id) from a V2 streaming
    source that implements both SupportsMetadataColumns and SupportsPushDownRequiredColumns,
    the analyzed plan expects the metadata column in the scan output, but Scan.readSchema()
    does not include it. Spark tries to read a column at an index the scan never produced.

  2. Root cause: pruneColumns is never called in streaming.
    In batch, V2ScanRelationPushDown calls SupportsPushDownRequiredColumns.pruneColumns
    with the required schema (which includes metadata columns resolved by AddMetadataColumns)
    before build(). In MicroBatchExecution.logicalPlan, the scan is built directly with
    table.newScanBuilder(options).build() — no pushdown of any kind is applied (a
    // TODO: operator pushdown comment marks this). Connectors that use pruneColumns to
    configure readSchema() — including whether to produce metadata columns — are never
    informed of what the query needs.

  3. This change fixes metadata column reads only, not column pruning.
    We call pruneColumns(output.toStructType) where output is the full analyzed relation
    output — all data columns plus any metadata columns added by AddMetadataColumns. This
    communicates required metadata columns to the scan builder so they appear in readSchema(),
    but does not prune data columns. Full column pruning in streaming, along with filter and
    aggregate pushdown, is deferred to the existing TODO.

Does this PR introduce any user-facing change?

Yes. Fixes a bug associated with metadata columns.

How was this patch tested?

Added a test in DataStreamTableAPISuite.

Was this patch authored or co-authored using generative AI tooling?

Yes

@zikangh zikangh changed the title [GH-56132][STREAMING] Call pruneColumns on V2 streaming scan builders in MicroBatchExecution [Spark-56132][STREAMING] Call pruneColumns on V2 streaming scan builders in MicroBatchExecution May 27, 2026
@zikangh zikangh changed the title [Spark-56132][STREAMING] Call pruneColumns on V2 streaming scan builders in MicroBatchExecution [SPARK-56132][STREAMING] Call pruneColumns on V2 streaming scan builders in MicroBatchExecution May 27, 2026

@gengliangwang gengliangwang left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

Problem: MicroBatchExecution.logicalPlan builds the V2 streaming scan via table.newScanBuilder(options).build() with no pushdown applied. The batch path (V2ScanRelationPushDown → PushDownUtils.pruneColumns) always invokes SupportsPushDownRequiredColumns.pruneColumns before build(), communicating the required schema including metadata columns added by AddMetadataColumns. Connectors that gate metadata-column exposure in Scan.readSchema() on the prune call work fine in batch but break in streaming — the analyzed plan's output contains the metadata column, the scan's readSchema() does not, and downstream binding crashes with ArrayIndexOutOfBoundsException.

Fix: Before build(), pattern-match the builder as SupportsPushDownRequiredColumns and call pruneColumns(output.toStructType). Passing the full analyzed output (data columns + metadata columns) tells the connector which metadata columns to expose in readSchema(), without actually pruning data columns. Full streaming pushdown (column pruning, filters, aggregates) is deferred to the existing // TODO: operator pushdown comment.

Key design decisions:

  • Inline match in logicalPlan rather than a streaming-equivalent of V2ScanRelationPushDown. Acceptable as a narrow fix because no streaming pushdown rule exists; the TODO comment marks this as future work.
  • Pass output.toStructType rather than a narrowed required schema — isolates the fix to metadata-column exposure, no data-column pruning side effect.
  • Skip the batch path's toOutputAttrs(scan.readSchema(), relation) rebind — trusts the connector to honor the prune schema. See the inline comment on line 231 for the trade-off.

Findings below are mostly suggestions / questions; the fix itself is sound and well-scoped.

val scanBuilder = table.newScanBuilder(options)
scanBuilder match {
case r: SupportsPushDownRequiredColumns =>
r.pruneColumns(output.toStructType)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ContinuousExecution.logicalPlan (sql/core/src/main/scala/org/apache/spark/sql/execution/streaming/continuous/ContinuousExecution.scala:95) has the identical pattern — table.newScanBuilder(options).build() with the same // TODO: operator pushdown comment. A query that selects a metadata column from a V2 continuous source implementing both SupportsMetadataColumns and SupportsPushDownRequiredColumns should hit the same ArrayIndexOutOfBoundsException. Either apply this fix symmetrically there, or note in the PR description why microbatch-only is sufficient.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch. done!

r.pruneColumns(output.toStructType)
case _ =>
}
val scan = scanBuilder.build()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The batch analogue in PushDownUtils.pruneColumns rebinds the relation's output from scan.readSchema() after build() (via toOutputAttrs(scan.readSchema(), relation)) so the relation reflects what the scan actually produces. Here we keep the analyzed output and trust the connector to produce a matching readSchema(). If a connector reorders fields or silently drops an unrecognized column, downstream binding will fail with a different cryptic error that looks like the bug this PR fixes — making future regressions confusing to diagnose. Consider either adopting the batch defensive rebind, or adding a short comment explaining why the analyzed output is safe to trust here.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not an issue here because we pass the entire output into pruneColumns.

try {
stream.addData(1, 2, 3)
q.processAllAvailable()
assert(recorded.called, "pruneColumns should have been called on the streaming scan builder")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This test asserts the pruneColumns API call happens, but doesn't actually exercise the user-visible bug. RecordingPruneScanBuilder.build() returns the unmodified inner MemoryStreamScanBuilder whose readSchema() is stream.fullSchema() = {value} — _seq is never produced by the scan regardless of what pruneColumns received. So the assertion passes even though the original ArrayIndexOutOfBoundsException scenario isn't reproduced here. The test protects against regression of the API call but not the failure mode described in the PR.

Consider adding a second test that uses a connector whose readSchema() actually depends on the prune input (conditionally exposing _seq when requested) and asserts via a memory sink (not noop) that selecting _seq returns expected values end-to-end. That would protect against connector-side breakage and document the contract the fix relies on.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added an E2E test.

}
}

test("pruneColumns called on SupportsPushDownRequiredColumns V2 streaming scan builder") {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Minor: adjacent tests use a SPARK-NNNNN: prefix (e.g. line 553 SPARK-44865:). Suggest matching that convention.

Suggested change
test("pruneColumns called on SupportsPushDownRequiredColumns V2 streaming scan builder") {
test("SPARK-56132: pruneColumns called on SupportsPushDownRequiredColumns V2 streaming scan builder") {

@zikangh
zikangh force-pushed the stack/prunecolumns-streaming branch from 5bbcfeb to 1e184cf Compare May 27, 2026 20:44
@zikangh zikangh changed the title [SPARK-56132][STREAMING] Call pruneColumns on V2 streaming scan builders in MicroBatchExecution [SPARK-56132][SS] Call pruneColumns on V2 streaming scan builders in MicroBatchExecution May 27, 2026
@zikangh zikangh changed the title [SPARK-56132][SS] Call pruneColumns on V2 streaming scan builders in MicroBatchExecution [SPARK-56132][SS] Call pruneColumns on V2 streaming to fix metadata reading issue May 27, 2026
@zikangh
zikangh requested a review from gengliangwang May 27, 2026 23:11
.booleanConf
.createWithDefault(true)

val STREAMING_V2_PRUNE_COLUMNS_ENABLED =

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

why do we need a new conf here? will there be regressions?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

removed.

"SupportsPushDownRequiredColumns. This communicates metadata columns to the scan " +
"builder so they appear in readSchema(). Disable if a connector's pruneColumns " +
"implementation is not compatible with being called in the streaming path.")
.withBindingPolicy(ConfigBindingPolicy.SESSION)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

why do we need this line?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

removed.

@zikangh
zikangh requested a review from gengliangwang May 28, 2026 17:44
@gengliangwang

Copy link
Copy Markdown
Member

@zikangh thanks for the work.
Since this is a bug fix, I am merging it to 4.2 as well. cc @huaxingao

gengliangwang pushed a commit that referenced this pull request May 28, 2026
…eading issue

### What changes were proposed in this pull request?

In `MicroBatchExecution.logicalPlan`, before calling `build()` on the V2 streaming scan
builder, call `SupportsPushDownRequiredColumns.pruneColumns(output.toStructType)` if the
builder supports it. `output` is the analyzed relation output, which already includes any
metadata columns the query references (added by the `AddMetadataColumns` rule).

### Why are the changes needed?

1. **Metadata column reads in V2 streaming crash with `ArrayIndexOutOfBoundsException`.**
   When a query selects a metadata column (e.g. `_metadata.row_id`) from a V2 streaming
   source that implements both `SupportsMetadataColumns` and `SupportsPushDownRequiredColumns`,
   the analyzed plan expects the metadata column in the scan output, but `Scan.readSchema()`
   does not include it. Spark tries to read a column at an index the scan never produced.

2. **Root cause: `pruneColumns` is never called in streaming.**
   In batch, `V2ScanRelationPushDown` calls `SupportsPushDownRequiredColumns.pruneColumns`
   with the required schema (which includes metadata columns resolved by `AddMetadataColumns`)
   before `build()`. In `MicroBatchExecution.logicalPlan`, the scan is built directly with
   `table.newScanBuilder(options).build()` — no pushdown of any kind is applied (a
   `// TODO: operator pushdown` comment marks this). Connectors that use `pruneColumns` to
   configure `readSchema()` — including whether to produce metadata columns — are never
   informed of what the query needs.

3. **This change fixes metadata column reads only, not column pruning.**
   We call `pruneColumns(output.toStructType)` where `output` is the full analyzed relation
   output — all data columns plus any metadata columns added by `AddMetadataColumns`. This
   communicates required metadata columns to the scan builder so they appear in `readSchema()`,
   but does not prune data columns. Full column pruning in streaming, along with filter and
   aggregate pushdown, is deferred to the existing TODO.

### Does this PR introduce _any_ user-facing change?

Yes. Fixes a bug associated with metadata columns.

### How was this patch tested?

Added a test in `DataStreamTableAPISuite`.

### Was this patch authored or co-authored using generative AI tooling?

Yes

Closes #56133 from zikangh/stack/prunecolumns-streaming.

Authored-by: Zikang Han <zikang.han@databricks.com>
Signed-off-by: Gengliang Wang <gengliang@apache.org>
(cherry picked from commit 2adc69a)
Signed-off-by: Gengliang Wang <gengliang@apache.org>
gengliangwang pushed a commit that referenced this pull request May 28, 2026
…eading issue

### What changes were proposed in this pull request?

In `MicroBatchExecution.logicalPlan`, before calling `build()` on the V2 streaming scan
builder, call `SupportsPushDownRequiredColumns.pruneColumns(output.toStructType)` if the
builder supports it. `output` is the analyzed relation output, which already includes any
metadata columns the query references (added by the `AddMetadataColumns` rule).

### Why are the changes needed?

1. **Metadata column reads in V2 streaming crash with `ArrayIndexOutOfBoundsException`.**
   When a query selects a metadata column (e.g. `_metadata.row_id`) from a V2 streaming
   source that implements both `SupportsMetadataColumns` and `SupportsPushDownRequiredColumns`,
   the analyzed plan expects the metadata column in the scan output, but `Scan.readSchema()`
   does not include it. Spark tries to read a column at an index the scan never produced.

2. **Root cause: `pruneColumns` is never called in streaming.**
   In batch, `V2ScanRelationPushDown` calls `SupportsPushDownRequiredColumns.pruneColumns`
   with the required schema (which includes metadata columns resolved by `AddMetadataColumns`)
   before `build()`. In `MicroBatchExecution.logicalPlan`, the scan is built directly with
   `table.newScanBuilder(options).build()` — no pushdown of any kind is applied (a
   `// TODO: operator pushdown` comment marks this). Connectors that use `pruneColumns` to
   configure `readSchema()` — including whether to produce metadata columns — are never
   informed of what the query needs.

3. **This change fixes metadata column reads only, not column pruning.**
   We call `pruneColumns(output.toStructType)` where `output` is the full analyzed relation
   output — all data columns plus any metadata columns added by `AddMetadataColumns`. This
   communicates required metadata columns to the scan builder so they appear in `readSchema()`,
   but does not prune data columns. Full column pruning in streaming, along with filter and
   aggregate pushdown, is deferred to the existing TODO.

### Does this PR introduce _any_ user-facing change?

Yes. Fixes a bug associated with metadata columns.

### How was this patch tested?

Added a test in `DataStreamTableAPISuite`.

### Was this patch authored or co-authored using generative AI tooling?

Yes

Closes #56133 from zikangh/stack/prunecolumns-streaming.

Authored-by: Zikang Han <zikang.han@databricks.com>
Signed-off-by: Gengliang Wang <gengliang@apache.org>
(cherry picked from commit 2adc69a)
Signed-off-by: Gengliang Wang <gengliang@apache.org>
murali-db pushed a commit to delta-io/delta that referenced this pull request Jun 5, 2026
…mapping combinations (#6910)

#### Which Delta project/connector is this regarding?

- [x] Spark

## Description

Adds streaming regression tests for row-tracking and CDC combination
scenarios, with DSv2 parity via \`V2ForceTest\`.

**Row-tracking metadata in streaming**
(\`DeltaSourceRowTrackingSuiteBase\` /
\`DeltaV2SourceRowTrackingSuite\`)

- V1 streaming (\`DeltaSource\` / \`StreamingRelation\`) does not expose
\`_metadata.row_id\`; the relation schema contains only user data
columns
- A CDC stream on a row-tracking table works normally when \`_metadata\`
is not projected — the Delta protocol prohibition fires only when
row-tracking fields are actually requested
- A CDC stream on a row-tracking + column-mapped table rejects
\`_metadata.row_id\` — column mapping must not bypass the protocol guard

**CDC streaming combinations** (added to \`DeltaCDCStreamSuiteBase\` /
\`DeltaV2CDCStreamSuite\`)

- A CDC stream on a partitioned table returns data and partition columns
correctly; CDC tail columns (\`_change_type\`, \`_commit_version\`,
etc.) are stripped
- A CDC stream on a column-mapped table returns correct logical column
values (id-mode and name-mode subclasses each exercise their own mode)

**V2 connector compatibility**

Each new base suite gets a \`V2ForceTest\` runner in \`spark-unified/\`.
Tests blocked on a pending Spark upstream fix (\`pruneColumns()\` in
\`MicroBatchExecution\`, apache/spark#56133) are classified in
\`shouldFailTests\` with a TODO rather than silently skipped.

## How was this patch tested?

Run locally against Spark 4.1 (\`sparkV2/testOnly\`):

| Suite | Passed | Ignored |
|---|---|---|
| \`DeltaSourceRowTrackingSuite\` (V1, 3 tests) | 3 | — |
| \`DeltaV2SourceRowTrackingSuite\` (V2) | 2 | 1 (pending
apache/spark#56133) |
| \`DeltaCDCStreamSuite\` (2 new tests) | 2 | — |
| \`DeltaV2CDCStreamSuite\` (2 new tests) | 2 | — |

## Does this PR introduce _any_ user-facing changes?

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants