laskoviymishka commented on code in PR #2106: URL: https://github.com/apache/iceberg-go/pull/2106#discussion_r4193432230
########## table/readtasks_residual_binding_internal_test.go: ########## @@ -0,0 +1,60 @@ +// 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" + "slices" + "testing" + + "github.com/apache/iceberg-go" + iceio "github.com/apache/iceberg-go/io" + "github.com/stretchr/testify/require" +) + +func TestReadTasksMixedResidualBindingPreservesInput(t *testing.T) { + 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) + scan := New(Identifier{"db", "tbl"}, metadata, "metadata.json", func(context.Context) (iceio.IO, error) { return iceio.NewMemFS(), nil }, nil).Scan() + unbound := iceberg.GreaterThan(iceberg.Reference("id"), int64(1)) + bound, err := iceberg.BindExpr(schema, unbound, true) + require.NoError(t, err) + for _, invalid := range []bool{false, true} { + t.Run(map[bool]string{false: "reusable mixed plan", true: "invalid bound residual after binding"}[invalid], func(t *testing.T) { + tasks := []FileScanTask{{Residual: bound}, {}, {Residual: unbound}, {Residual: bound}} + if invalid { + 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) + tasks[len(tasks)-1].Residual = wrongBound + } + before := slices.Clone(tasks) + for range 2 { + _, _, err := scan.ReadTasks(t.Context(), tasks) + if invalid { + require.ErrorIs(t, err, iceberg.ErrInvalidArgument) + require.ErrorContains(t, err, "field ID 2") + } else { + require.NoError(t, err) + } + require.Equal(t, before, tasks) Review Comment: This only asserts the input is untouched, which misses the regression that'd actually bite: nothing observes the slice handed to `GetRecords`. A version that cloned but forgot `readTasks[i].Residual = boundResidual`, or that reverted to an unconditional `slices.Clone`, passes every case here. I'd pull the bind loop into a helper returning `([]FileScanTask, error)` and assert on the output: bound residuals at the right indices, the all-bound/nil case aliasing the input (`&out[0] == &tasks[0]`), and a changed case getting a fresh backing array. A first-task-unbound case (clone at index 0) and an all-unbound case would also close the COW-boundary gaps the current `{bound, nil, unbound, bound}` fixture leaves. Minor on the same assertion: `require.Equal(t, before, tasks)` falls back to `reflect.DeepEqual`, so it won't distinguish the unbound slot being re-bound to an equal-but-distinct expression from being left untouched. `require.Same` on `tasks[2].Residual` pins the invariant more precisely. ########## table/scanner.go: ########## @@ -2202,17 +2202,26 @@ func (scan *Scan) ReadTasks(ctx context.Context, tasks []FileScanTask) (*arrow.S // Bind task residuals against the schema selected by this scan, which may // be an older snapshot schema rather than the table's current schema. Keep // the caller's task slice untouched because the same plan may be reused. - readTasks := slices.Clone(tasks) - for i := range readTasks { - if readTasks[i].Residual == nil { + readTasks := tasks Review Comment: On the no-change path `readTasks` is the caller's backing array now, not a copy, and we hand it straight to `GetRecords`, which returns a lazy `iter.Seq2` that ranges over the slice well after `ReadTasks` has returned. The old unconditional `slices.Clone` snapshotted the plan; without it, a caller that reuses, re-slices, sorts, or overwrites its `[]FileScanTask` while an earlier iterator is still draining races or reads the wrong tasks. It's latent today because `GetRecords` only reads the elements, but nothing documents or enforces that, and the retained comment ("Keep the caller's task slice untouched") only covers our half of it, not caller mutation during iteration. I'd either keep the defensive clone, or make the new contract explicit: a godoc line on `ReadTasks` that the slice must stay unmodified until the returned iterator is exhausted or abandoned, plus a note at the alias site that `readTasks` aliases `tasks` so everything downstream must treat it read-only. If we keep the alias, I'd pin "downstream only reads" with a test so a future write into a task element gets caught rather than silently corrupting a reused plan. The clone is a single slice copy the benchmark shows is cheap next to reading the files, so dropping it in favor of an undocumented invariant feels like the riskier trade. ########## table/readtasks_residual_binding_internal_test.go: ########## @@ -0,0 +1,60 @@ +// 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" + "slices" + "testing" + + "github.com/apache/iceberg-go" + iceio "github.com/apache/iceberg-go/io" + "github.com/stretchr/testify/require" +) + +func TestReadTasksMixedResidualBindingPreservesInput(t *testing.T) { + 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) + scan := New(Identifier{"db", "tbl"}, metadata, "metadata.json", func(context.Context) (iceio.IO, error) { return iceio.NewMemFS(), nil }, nil).Scan() + unbound := iceberg.GreaterThan(iceberg.Reference("id"), int64(1)) + bound, err := iceberg.BindExpr(schema, unbound, true) + require.NoError(t, err) + for _, invalid := range []bool{false, true} { + t.Run(map[bool]string{false: "reusable mixed plan", true: "invalid bound residual after binding"}[invalid], func(t *testing.T) { Review Comment: The `map[bool]string{...}[invalid]` subtest-name trick builds a map every iteration just to pick a label, and the `for range 2` loop leaves the reader to infer that the point is proving the plan is reusable. A small table of `{name, wantErr}` cases with `t.Run(tc.name, ...)` reads cleaner, and a `// run twice: the plan must be reusable` on the loop (or a named loop var surfaced in the failure message) spells out the intent. ########## table/arrow_scanner.go: ########## @@ -1689,29 +1689,34 @@ func validateBoundFilter(schema *iceberg.Schema, filter iceberg.BooleanExpressio return validationErr } -func bindTaskFilter(schema *iceberg.Schema, filter iceberg.BooleanExpression, caseSensitive bool) (iceberg.BooleanExpression, error) { +func bindTaskFilter(schema *iceberg.Schema, filter iceberg.BooleanExpression, caseSensitive bool) (iceberg.BooleanExpression, bool, error) { Review Comment: Since `changed` is now load-bearing for the clone decision, I'd add a godoc line spelling out the contract: it's true only when a new expression was produced (the unbound-rebind path), and when it's false the returned filter is the exact input value. Three of the four callers discard the flag, so the one that relies on it reads better with the invariant written down. Worth a line on `FileScanTask.Residual` too that callers may supply either a bound or unbound expression. ########## table/arrow_scanner_bench_test.go: ########## @@ -687,3 +688,61 @@ func benchmarkTaskResidualJSON(rowCount int) string { return result.String() } + +var benchmarkReadTasksCount int + +func BenchmarkArrowScanReadTasksResidualBinding(b *testing.B) { + schema := iceberg.NewSchema(1, + iceberg.NestedField{ID: 1, Name: "id", Type: iceberg.PrimitiveTypes.Int64, Required: true}, + ) + metadata, err := NewMetadata( + schema, iceberg.UnpartitionedSpec, UnsortedSortOrder, "mem://benchmark/read-tasks", nil, + ) + if err != nil { + b.Fatal(err) + } + fs := iceio.NewMemFS() + scan := New(Identifier{"benchmark", "read_tasks"}, metadata, "metadata.json", + func(context.Context) (iceio.IO, error) { return fs, nil }, nil).Scan() + boundResidual, err := iceberg.BindExpr(schema, + iceberg.GreaterThan(iceberg.Reference("id"), int64(1)), true) + if err != nil { + b.Fatal(err) + } + unboundResidual := iceberg.GreaterThan(iceberg.Reference("id"), int64(1)) + + for _, workload := range []struct { + taskCount int + name string + residual iceberg.BooleanExpression + }{ + {taskCount: 1_000, name: "bound", residual: boundResidual}, + {taskCount: 10_000, name: "bound", residual: boundResidual}, + {taskCount: 100_000, name: "bound", residual: boundResidual}, + {taskCount: 100_000, name: "nil"}, + {taskCount: 100_000, name: "unbound", residual: unboundResidual}, + {taskCount: 100_000, name: "late_unbound", residual: boundResidual}, + } { + b.Run(fmt.Sprintf("tasks=%d/residual=%s", workload.taskCount, workload.name), func(b *testing.B) { + tasks := make([]FileScanTask, workload.taskCount) + for i := range tasks { + tasks[i].Residual = workload.residual + } + if workload.name == "late_unbound" { + tasks[len(tasks)-1].Residual = unboundResidual + } + + b.ReportAllocs() + b.ResetTimer() + for b.Loop() { + _, records, err := scan.ReadTasks(b.Context(), tasks) + if err != nil { + b.Fatal(err) + } + benchmarkReadTasksCount = len(tasks) Review Comment: `benchmarkReadTasksCount = len(tasks)` sinks a loop-invariant that was never at risk of being elided; the result that matters is `records`, which `runtime.KeepAlive` already pins, so the package global and the `runtime` import can both go. The iterator is never drained, so this measures setup plus the bind loop only. That's fine for the allocation claim (allocs/op and B/op are the real signal, the wall-time figures are end-to-end), but a one-line comment that iteration is skipped on purpose would save the next reader the puzzle. -- 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]
