Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 16 additions & 1 deletion table/arrow_scanner.go
Original file line number Diff line number Diff line change
Expand Up @@ -869,7 +869,10 @@ func (as *arrowScan) projectedFieldIDs(rowFilter iceberg.BooleanExpression, equa
}

func (as *arrowScan) scanInvariants(tableProperties iceberg.Properties) (*arrowScanInvariants, error) {
projectedIDs, err := as.projectedFieldIDs(as.boundRowFilter, nil)
// Filter columns are added per task below. A local task may have a
// partition-elided residual that no longer references the original filter
// fields.
projectedIDs, err := as.projectedFieldIDs(nil, nil)
if err != nil {
return nil, err
}
Expand All @@ -883,6 +886,18 @@ func (as *arrowScan) scanInvariants(tableProperties iceberg.Properties) (*arrowS
}

func (as *arrowScan) addTaskProjectedFieldIDs(invariants *arrowScanInvariants, tasks []FileScanTask) error {
// Tasks without a residual still use the scan's original filter. Add those
// fields once, then add only the actual residual fields for other tasks.
for _, task := range tasks {
if task.Residual == nil {
if err := addFilterFieldIDs(invariants.projectedIDs, as.boundRowFilter); err != nil {
return err
}

break
}
}

for _, task := range tasks {
if task.Residual == nil {
continue
Expand Down
37 changes: 34 additions & 3 deletions table/arrow_scanner_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -149,7 +149,7 @@ func TestArrowScanSnapshotsInvariants(t *testing.T) {
assert.Equal(t, 1, metadata.nameMappingCalls)
}

func TestArrowScanAddTaskProjectedFieldIDsSkipsNilResiduals(t *testing.T) {
func TestArrowScanAddTaskProjectedFieldIDsUsesTaskResiduals(t *testing.T) {
schema := iceberg.NewSchema(1,
iceberg.NestedField{ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64},
iceberg.NestedField{ID: 2, Name: "data", Type: iceberg.PrimitiveTypes.String},
Expand All @@ -162,9 +162,9 @@ func TestArrowScanAddTaskProjectedFieldIDsSkipsNilResiduals(t *testing.T) {
scanSchema: schema,
boundRowFilter: boundFilter,
}
invariants := &arrowScanInvariants{projectedIDs: set[int]{1: {}}}
invariants := &arrowScanInvariants{projectedIDs: set[int]{}}
tasks := []FileScanTask{
{},
{Residual: iceberg.AlwaysTrue{}},
{Residual: iceberg.EqualTo(iceberg.Reference("data"), "value")},
{},
}
Expand All @@ -174,6 +174,37 @@ func TestArrowScanAddTaskProjectedFieldIDsSkipsNilResiduals(t *testing.T) {
assert.Equal(t, set[int]{1: {}, 2: {}}, invariants.projectedIDs)
}

func TestArrowScanInvariantsStartWithRequestedFields(t *testing.T) {
schema := iceberg.NewSchema(1,
iceberg.NestedField{ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64},
iceberg.NestedField{ID: 2, Name: "data", Type: iceberg.PrimitiveTypes.String},
)
projected, err := schema.Select(true, "data")
require.NoError(t, err)
metadata, err := NewMetadata(
schema, iceberg.UnpartitionedSpec, UnsortedSortOrder, "mem://test/table", iceberg.Properties{},
)
require.NoError(t, err)

boundFilter, err := iceberg.BindExpr(schema,
iceberg.EqualTo(iceberg.Reference("id"), int64(1)), true)
require.NoError(t, err)

scanner := &arrowScan{
metadata: metadata,
scanSchema: schema,
projectedSchema: projected,
boundRowFilter: boundFilter,
}
invariants, err := scanner.scanInvariants(metadata.Properties())
require.NoError(t, err)
assert.Equal(t, set[int]{2: {}}, invariants.projectedIDs)

err = scanner.addTaskProjectedFieldIDs(invariants, []FileScanTask{{Residual: iceberg.AlwaysTrue{}}})
require.NoError(t, err)
assert.Equal(t, set[int]{2: {}}, invariants.projectedIDs)
}

func TestEnrichRecordsWithPosDeleteFields(t *testing.T) {
testSchema := arrow.NewSchema([]arrow.Field{
{Name: "first_name", Type: &arrow.StringType{}, Nullable: false},
Expand Down
Loading
Loading