laskoviymishka commented on code in PR #2106:
URL: https://github.com/apache/iceberg-go/pull/2106#discussion_r4198924417
##########
table/scanner.go:
##########
@@ -2161,6 +2193,11 @@ func (scan *Scan) ToArrowRecords(ctx context.Context)
(*arrow.Schema, iter.Seq2[
// reached; if no such task is processed, the file is not read and its error
is not
// returned. The returned iterator is single-use.
//
+// The caller must not modify tasks or any task element until the returned
+// iterator is exhausted or abandoned. When no residual needs binding (each is
Review Comment:
The thing I'd tighten is "abandoned". It's not an event a caller can
observe: if the range loop is never entered or breaks early, nothing marks when
the hold ends, and with `WithMaxConcurrency` above one workers may still be
draining the aliased slice after a break. I'd tie it to the loop, safe to
mutate once the range over the returned iterator has returned.
Two words on the rest of the wording. "may retain the backing array"
undersells it, since the nested slices (`File`, `DeleteFiles`,
`EqualityDeleteFiles`) are shared regardless (the old clone was shallow too),
so "treat tasks as immutable for the iterator's lifetime" is the honest
framing. And because the alias only triggers on the all-bound/nil shape, a
caller who reuses the plan buffer sees corruption on one plan shape and nothing
on another, so stating it absolutely is what keeps that from being a latent
trap.
##########
table/arrow_scanner.go:
##########
@@ -1724,7 +1732,9 @@ func (as *arrowScan) rowFilterForTask(task FileScanTask)
(iceberg.BooleanExpress
filterSchema = as.filterSchema
}
- return bindTaskFilter(filterSchema, task.Residual, as.caseSensitive)
+ bound, _, err := bindTaskFilter(filterSchema, task.Residual,
as.caseSensitive)
Review Comment:
Not blocking, but worth a note since this is a perf PR: on the all-bound
path we now validate each residual twice, once in `bindReadTasksResiduals` and
once here per task, so we've dropped the allocation but kept the
`VisitExpr`/`validateBoundFilter` cost. The benchmark doesn't drain the
iterator, so that CPU stays off the headline number. Either skip validation
here when the helper already validated, or call it out in the PR description so
the 13.6MB to 4KB figure isn't read as the whole story.
##########
table/scanner.go:
##########
@@ -2115,7 +2115,9 @@ type FileScanTask struct {
Start, Length int64
// Residual is the portion of the scan filter that must still be
evaluated
// for this task. Local and remote planners may simplify the original
filter using
- // file metadata; nil means the caller did not provide a task residual.
+ // file metadata; nil means the caller did not provide a task residual.
Callers
+ // may supply either a bound or unbound expression; ReadTasks validates
bound
Review Comment:
Worth one more clause here: a residual that mixes bound and unbound
predicates is rejected (`bindTaskFilter`'s `hasBound && hasUnbound` branch),
not bound. As written, "either a bound or unbound expression" reads as if mixed
is fine too.
##########
table/readtasks_residual_binding_internal_test.go:
##########
@@ -0,0 +1,234 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package table
+
+import (
+ "context"
+ "os"
+ "path/filepath"
+ "testing"
+
+ "github.com/apache/iceberg-go"
+ iceio "github.com/apache/iceberg-go/io"
+ "github.com/stretchr/testify/require"
+)
+
+func residualBindingTestScan(t *testing.T) (*Scan, *iceberg.Schema) {
+ t.Helper()
+
+ schema := iceberg.NewSchema(0,
+ iceberg.NestedField{ID: 1, Name: "id", Type:
iceberg.PrimitiveTypes.Int64},
+ )
+ metadata, err := NewMetadata(
+ schema, iceberg.UnpartitionedSpec, UnsortedSortOrder,
"mem://mixed-residuals", nil,
+ )
+ require.NoError(t, err)
+ memFS := iceio.NewMemFS()
+ tbl := New(
+ Identifier{"db", "tbl"}, metadata, "metadata.json",
+ func(context.Context) (iceio.IO, error) { return memFS, nil },
nil,
+ )
+
+ return tbl.Scan(), schema
+}
+
+func TestBindReadTasksResidualsCopyOnWrite(t *testing.T) {
+ _, schema := residualBindingTestScan(t)
+ unbound := iceberg.GreaterThan(iceberg.Reference("id"), int64(1))
+ bound, err := iceberg.BindExpr(schema, unbound, true)
+ require.NoError(t, err)
+
+ tests := []struct {
+ name string
+ residuals []iceberg.BooleanExpression
+ wantAlias bool
+ }{
+ {name: "all bound and nil", residuals:
[]iceberg.BooleanExpression{bound, nil, bound}, wantAlias: true},
+ {name: "mixed", residuals: []iceberg.BooleanExpression{bound,
nil, unbound, bound}},
+ {name: "first task unbound", residuals:
[]iceberg.BooleanExpression{unbound, bound}},
+ {name: "all unbound", residuals:
[]iceberg.BooleanExpression{unbound, unbound}},
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ tasks := make([]FileScanTask, len(tt.residuals))
+ for i, residual := range tt.residuals {
+ tasks[i].Residual = residual
+ }
+
+ got, err := bindReadTasksResiduals(schema, tasks, true)
+ require.NoError(t, err)
+ require.Len(t, got, len(tasks))
+ if len(tasks) > 0 {
+ require.Equal(t, tt.wantAlias, &got[0] ==
&tasks[0])
+ }
+
+ for i, original := range tt.residuals {
+ if original == nil {
+ require.Nil(t, got[i].Residual)
+
+ continue
+ }
+
+ // The input plan is never rewritten, even when
the output needs binding.
+ require.Same(t, original, tasks[i].Residual)
+ state, visitErr :=
iceberg.VisitExpr(got[i].Residual, filterBindingVisitor{})
+ require.NoError(t, visitErr)
+ require.True(t, state.hasBound)
Review Comment:
This conditional makes the identity check fail open: if `VisitExpr` errors
it's silently skipped, and it re-derives the same bound/unbound split the
production code uses, so the assertion is partly checking the code against
itself. The table already knows which cases are bound, so I'd put a `wantSame
[]bool` (or a `preBound` flag) on each case and assert `require.Same` /
`require.NotSame` directly.
##########
table/readtasks_residual_binding_internal_test.go:
##########
@@ -0,0 +1,234 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package table
+
+import (
+ "context"
+ "os"
+ "path/filepath"
+ "testing"
+
+ "github.com/apache/iceberg-go"
+ iceio "github.com/apache/iceberg-go/io"
+ "github.com/stretchr/testify/require"
+)
+
+func residualBindingTestScan(t *testing.T) (*Scan, *iceberg.Schema) {
+ t.Helper()
+
+ schema := iceberg.NewSchema(0,
+ iceberg.NestedField{ID: 1, Name: "id", Type:
iceberg.PrimitiveTypes.Int64},
+ )
+ metadata, err := NewMetadata(
+ schema, iceberg.UnpartitionedSpec, UnsortedSortOrder,
"mem://mixed-residuals", nil,
+ )
+ require.NoError(t, err)
+ memFS := iceio.NewMemFS()
+ tbl := New(
+ Identifier{"db", "tbl"}, metadata, "metadata.json",
+ func(context.Context) (iceio.IO, error) { return memFS, nil },
nil,
+ )
+
+ return tbl.Scan(), schema
+}
+
+func TestBindReadTasksResidualsCopyOnWrite(t *testing.T) {
+ _, schema := residualBindingTestScan(t)
+ unbound := iceberg.GreaterThan(iceberg.Reference("id"), int64(1))
+ bound, err := iceberg.BindExpr(schema, unbound, true)
+ require.NoError(t, err)
+
+ tests := []struct {
+ name string
+ residuals []iceberg.BooleanExpression
+ wantAlias bool
+ }{
+ {name: "all bound and nil", residuals:
[]iceberg.BooleanExpression{bound, nil, bound}, wantAlias: true},
+ {name: "mixed", residuals: []iceberg.BooleanExpression{bound,
nil, unbound, bound}},
+ {name: "first task unbound", residuals:
[]iceberg.BooleanExpression{unbound, bound}},
+ {name: "all unbound", residuals:
[]iceberg.BooleanExpression{unbound, unbound}},
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ tasks := make([]FileScanTask, len(tt.residuals))
+ for i, residual := range tt.residuals {
+ tasks[i].Residual = residual
+ }
+
+ got, err := bindReadTasksResiduals(schema, tasks, true)
+ require.NoError(t, err)
+ require.Len(t, got, len(tasks))
+ if len(tasks) > 0 {
+ require.Equal(t, tt.wantAlias, &got[0] ==
&tasks[0])
+ }
+
+ for i, original := range tt.residuals {
+ if original == nil {
+ require.Nil(t, got[i].Residual)
+
+ continue
+ }
+
+ // The input plan is never rewritten, even when
the output needs binding.
+ require.Same(t, original, tasks[i].Residual)
+ state, visitErr :=
iceberg.VisitExpr(got[i].Residual, filterBindingVisitor{})
+ require.NoError(t, visitErr)
+ require.True(t, state.hasBound)
+ require.False(t, state.hasUnbound)
+ if stateBefore, stateErr :=
iceberg.VisitExpr(original, filterBindingVisitor{}); stateErr == nil &&
!stateBefore.hasUnbound {
+ require.Same(t, original,
got[i].Residual)
+ }
+ }
+ })
+ }
+}
+
+func TestReadTasksResidualPlanIsReusable(t *testing.T) {
+ scan, schema := residualBindingTestScan(t)
+ unbound := iceberg.GreaterThan(iceberg.Reference("id"), int64(1))
+ bound, err := iceberg.BindExpr(schema, unbound, true)
+ require.NoError(t, err)
+
+ wrongSchema := iceberg.NewSchema(0,
+ iceberg.NestedField{ID: 2, Name: "other", Type:
iceberg.PrimitiveTypes.Int64},
+ )
+ wrongBound, err := iceberg.BindExpr(
+ wrongSchema, iceberg.EqualTo(iceberg.Reference("other"),
int64(1)), true,
+ )
+ require.NoError(t, err)
+
+ tests := []struct {
+ name string
+ tasks []FileScanTask
+ wantErr bool
+ }{
+ {
+ name: "mixed plan",
+ tasks: []FileScanTask{{Residual: bound}, {}, {Residual:
unbound}, {Residual: bound}},
+ },
+ {
+ name: "invalid bound residual",
+ tasks: []FileScanTask{{Residual: bound}, {Residual:
unbound}, {Residual: wrongBound}},
+ wantErr: true,
+ },
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ originals := make([]iceberg.BooleanExpression,
len(tt.tasks))
+ for i := range tt.tasks {
+ originals[i] = tt.tasks[i].Residual
+ }
+
+ // Run twice to prove the caller-owned plan remains
reusable.
+ for range 2 {
+ _, _, err := scan.ReadTasks(t.Context(),
tt.tasks)
+ if tt.wantErr {
+ require.ErrorIs(t, err,
iceberg.ErrInvalidArgument)
+ require.ErrorContains(t, err, "field ID
2")
+ } else {
+ require.NoError(t, err)
+ }
+ for i, original := range originals {
+ if original == nil {
+ require.Nil(t,
tt.tasks[i].Residual)
+ } else {
+ require.Same(t, original,
tt.tasks[i].Residual)
+ }
+ }
+ }
+ })
+ }
+}
+
+func TestReadTasksAlreadyBoundTasksRemainReadOnly(t *testing.T) {
+ schema := iceberg.NewSchema(0,
+ iceberg.NestedField{ID: 1, Name: "id", Type:
iceberg.PrimitiveTypes.Int64},
+ )
+ location := t.TempDir()
+ metadata, err := NewMetadata(
+ schema, iceberg.UnpartitionedSpec, UnsortedSortOrder, location,
nil,
+ )
+ require.NoError(t, err)
+ tbl := New(
+ Identifier{"db", "tbl"}, metadata, filepath.Join(location,
"metadata.json"),
+ func(context.Context) (iceio.IO, error) { return
iceio.LocalFS{}, nil }, nil,
+ )
+ scan := tbl.Scan(WithMaxConcurrency(4))
+
+ unbound := iceberg.GreaterThan(iceberg.Reference("id"), int64(1))
+ bound, err := iceberg.BindExpr(schema, unbound, true)
+ require.NoError(t, err)
+ arrowSchema, err := SchemaToArrowSchema(schema, nil, false, false)
+ require.NoError(t, err)
+
+ paths := []string{
+ filepath.Join(location, "data-1.parquet"),
+ filepath.Join(location, "data-2.parquet"),
+ }
+ rows := []string{
+ `[{"id":2}]`,
+ `[{"id":3}]`,
+ }
+ tasks := make([]FileScanTask, len(paths))
+ for i, path := range paths {
+ writeParquetFile(t, path, arrowSchema, rows[i])
Review Comment:
This is the round-1 ask that's still open. These `require.Same` checks only
prove the caller's input wasn't mutated, which the old unconditional clone
satisfied just as well, so nothing here observes the slice `GetRecords`
actually receives. If `ReadTasks` reverted to cloning unconditionally, or
passed `tasks` instead of `readTasks`, this test stays green, and the
"ReadOnly" in the name overstates what it checks.
The hazard lives at the `ReadTasks` to `GetRecords` boundary, so that's
where I'd pin it: a seam (a package-private reader hook or func var) that lets
us assert on the `[]FileScanTask` `GetRecords` gets, `&received[0] ==
&tasks[0]` for all-bound, a clone with bound residuals for mixed. Reading real
parquet with an unbound residual and asserting the filtered result would also
do it. If neither lands this round, I'd at least rename to say what it checks
(input preserved) and drop ReadOnly.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]