laskoviymishka commented on code in PR #2028:
URL: https://github.com/apache/iceberg-go/pull/2028#discussion_r4064203620
##########
table/arrow_filter_plan.go:
##########
@@ -27,6 +27,7 @@ import (
"github.com/apache/arrow-go/v18/arrow/compute/exprs"
"github.com/apache/arrow-go/v18/parquet/metadata"
"github.com/apache/iceberg-go"
+ iceberginternal "github.com/apache/iceberg-go/internal"
Review Comment:
small thing while we're here — this pulls in the internal package as
`iceberginternal`, but `arrow_scanner.go` right next door already aliases the
same import as `iceinternal`. Two names for one package inside the same package
reads a bit odd. I'd match the existing `iceinternal` to keep it consistent.
##########
table/arrow_filter_plan.go:
##########
@@ -159,52 +160,24 @@ func (as *arrowScan) cachedFileFilterPlans(fileSchema
*iceberg.Schema, includePr
return plans, nil
}
-func physicalSchemaKey(fileSchema *iceberg.Schema) (string, error) {
+func physicalSchemaKey(fileSchema *iceberg.Schema) (key string, err error) {
+ defer func() {
+ if r := recover(); r != nil {
Review Comment:
I'd tighten this recover before merge. As written it catches every panic
from the whole call chain, not just the nil/invalid-type cases you're guarding
— so a future off-by-one or missed nil-check in `writePhysicalTypeKey` (say
when someone adds a new container case) gets silently rewritten into
`ErrInvalidSchema` on a perfectly valid schema, and the stack trace that
would've pointed at the real bug is gone.
Two ways out: guard the invalid inputs explicitly (a `typ == nil` check at
the top of `writePhysicalTypeKey`, nil-pointer checks before field access) and
drop the blanket recover, or keep the escape hatch but make it own its panics —
wrap the intended panic in a sentinel type and re-panic on anything else:
```go
if _, ok := r.(physicalKeyPanic); !ok {
panic(r) // not ours
}
```
Either is fine, but I'd rather we not swallow unrelated panics here. wdyt?
##########
table/arrow_filter_plan.go:
##########
@@ -246,11 +221,52 @@ func writePhysicalTypeKey(builder *strings.Builder, typ
iceberg.Type) {
builder.WriteByte('0')
}
writePhysicalTypeKey(builder, t.ValueType)
- case nil:
- builder.WriteString("n")
- default:
+ case iceberg.BooleanType:
+ builder.WriteByte('b')
+ case iceberg.Int32Type:
+ builder.WriteByte('i')
+ case iceberg.Int64Type:
+ builder.WriteByte('j')
+ case iceberg.Float32Type:
+ builder.WriteByte('k')
+ case iceberg.Float64Type:
+ builder.WriteByte('d')
+ case iceberg.DateType:
+ builder.WriteByte('D')
+ case iceberg.TimeType:
+ builder.WriteByte('T')
+ case iceberg.TimestampType:
+ builder.WriteByte('t')
+ case iceberg.TimestampTzType:
+ builder.WriteByte('z')
+ case iceberg.StringType:
+ builder.WriteByte('S')
+ case iceberg.UUIDType:
+ builder.WriteByte('u')
+ case iceberg.BinaryType:
+ builder.WriteByte('B')
+ case iceberg.TimestampNsType:
+ builder.WriteByte('n')
+ case iceberg.TimestampTzNsType:
+ builder.WriteByte('N')
+ case iceberg.UnknownType:
+ builder.WriteByte('U')
+ case iceberg.VariantType:
+ builder.WriteByte('v')
+ case iceberg.FixedType:
+ builder.WriteByte('f')
Review Comment:
not a correctness issue — the length-prefixed field names keep this
collision-free, so an `f8:` in the type slot can't be confused with a field
record. But `'f'` now means both "start of a field" and "FixedType" depending
on position, which makes a key like `f1:1:a:0f8:` a little hard to read by eye.
I'd give FixedType its own byte (`'F'` or `'x'`) just to keep the two meanings
separate. Minor, take it or leave it.
##########
table/arrow_filter_plan.go:
##########
@@ -159,52 +160,24 @@ func (as *arrowScan) cachedFileFilterPlans(fileSchema
*iceberg.Schema, includePr
return plans, nil
}
-func physicalSchemaKey(fileSchema *iceberg.Schema) (string, error) {
+func physicalSchemaKey(fileSchema *iceberg.Schema) (key string, err error) {
+ defer func() {
+ if r := recover(); r != nil {
+ key = ""
Review Comment:
`key` is a named return, so it's already `""` here — the panic happens
before `return builder.String(), nil` ever runs, so nothing could have set it.
I'd drop the assignment; as-is it reads like `key` might have been partially
written and we're resetting it, which isn't the case.
##########
table/arrow_filter_plan_test.go:
##########
@@ -196,3 +196,138 @@ func
TestCompiledFileFilterPlansDisablePruningForMissingInitialDefault(t *testin
assert.True(t, plans.pruning.statsFilter.Equals(iceberg.AlwaysTrue{}))
assert.Empty(t, plans.pruning.bloomPreds)
}
+
+func TestPhysicalSchemaKeyIgnoresMetadata(t *testing.T) {
+ fields := []iceberg.NestedField{{
+ ID: 1, Name: "root", Type: &iceberg.StructType{FieldList:
[]iceberg.NestedField{
+ {ID: 2, Name: "value", Type:
iceberg.PrimitiveTypes.Int64, Required: true},
+ }},
+ }}
+ original := iceberg.NewSchema(1, fields...)
+ key, err := physicalSchemaKey(original)
+ require.NoError(t, err)
+
+ fields = original.Fields()
+ fields[0].Doc = "root documentation"
+ child := &fields[0].Type.(*iceberg.StructType).FieldList[0]
+ child.Doc = "value documentation"
+ child.InitialDefault = int64(7)
+ child.WriteDefault = int64(9)
+ withMetadata := iceberg.NewSchemaWithIdentifiers(2, []int{2}, fields...)
+ metadataKey, err := physicalSchemaKey(withMetadata)
+ require.NoError(t, err)
+ assert.Equal(t, key, metadataKey)
+
+ emptyKey, err := physicalSchemaKey(iceberg.NewSchema(1))
+ require.NoError(t, err)
+ assert.Empty(t, emptyKey)
+}
+
+func TestPhysicalSchemaKeyDistinguishesNestedLayouts(t *testing.T) {
Review Comment:
this covers the structural mutations nicely, but the new bit that's easiest
to get wrong is the single-byte tag table itself — 14-ish tags assigned by
hand. A copy-paste dupe (the timestamp variants `t/z/n/N`, or `'n'` which was
the old nil tag now reused for TimestampNs) would silently alias two schemas
onto one cache key and apply the wrong compiled plan, and nothing here would
catch it.
Could we add a table-driven test that builds a one-field schema per
primitive type and asserts every pair produces a distinct key? That pins down
the property this whole change rests on.
##########
table/arrow_filter_plan.go:
##########
@@ -246,11 +221,52 @@ func writePhysicalTypeKey(builder *strings.Builder, typ
iceberg.Type) {
builder.WriteByte('0')
}
writePhysicalTypeKey(builder, t.ValueType)
- case nil:
- builder.WriteString("n")
- default:
+ case iceberg.BooleanType:
+ builder.WriteByte('b')
+ case iceberg.Int32Type:
+ builder.WriteByte('i')
+ case iceberg.Int64Type:
+ builder.WriteByte('j')
+ case iceberg.Float32Type:
+ builder.WriteByte('k')
+ case iceberg.Float64Type:
+ builder.WriteByte('d')
+ case iceberg.DateType:
+ builder.WriteByte('D')
+ case iceberg.TimeType:
+ builder.WriteByte('T')
+ case iceberg.TimestampType:
+ builder.WriteByte('t')
+ case iceberg.TimestampTzType:
+ builder.WriteByte('z')
+ case iceberg.StringType:
+ builder.WriteByte('S')
+ case iceberg.UUIDType:
+ builder.WriteByte('u')
+ case iceberg.BinaryType:
+ builder.WriteByte('B')
+ case iceberg.TimestampNsType:
+ builder.WriteByte('n')
+ case iceberg.TimestampTzNsType:
+ builder.WriteByte('N')
+ case iceberg.UnknownType:
+ builder.WriteByte('U')
+ case iceberg.VariantType:
+ builder.WriteByte('v')
+ case iceberg.FixedType:
+ builder.WriteByte('f')
+ writePhysicalInt(builder, t.Len())
+ case iceberg.DecimalType:
+ builder.WriteByte('q')
+ writePhysicalInt(builder, t.Precision())
+ writePhysicalInt(builder, t.Scale())
+ case iceberg.PrimitiveType:
+ // Keep custom and parameterized primitive types (for example
geometry)
+ // compatible with the previous structural key encoding.
builder.WriteByte('p')
writePhysicalString(builder, typ.String())
+ default:
+ panic(fmt.Errorf("unsupported physical type: %T", typ))
Review Comment:
the panic value is an `error` built with `fmt.Errorf`, but recover just
renders it with `%v`, so it's constructed as an error and never used as one
(the `%w` wrap in the outer message is for `ErrInvalidSchema`, not this).
`panic(fmt.Sprintf(...))` gives identical output and reads more clearly as a
one-way control-flow escape.
##########
table/arrow_scanner_bench_test.go:
##########
@@ -405,6 +405,97 @@ func BenchmarkArrowScanFilterPlanSetup(b *testing.B) {
})
}
+func BenchmarkArrowScanFilterPlanCacheHit(b *testing.B) {
+ for _, tc := range benchmarkPhysicalSchemaCases() {
+ b.Run(tc.name, func(b *testing.B) {
+ scan := &arrowScan{boundRowFilter:
iceberg.AlwaysTrue{}, caseSensitive: true}
+ if _, err := scan.cachedFileFilterPlans(tc.schema,
true); err != nil {
+ b.Fatal(err)
+ }
+
+ b.ReportAllocs()
Review Comment:
the sibling benches (the cache_hit sub-bench in
`BenchmarkArrowScanFilterPlanSetup`) do `b.ResetTimer()` after warm-up before
the loop, and this one skips it. In practice `b.Loop()` manages the timer
itself so the warm-up call probably isn't distorting the number, but it's
inconsistent with the rest of the file. I'd add the `ResetTimer()` after the
warm-up `cachedFileFilterPlans` just to match the idiom.
--
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]