diff --git a/weed/worker/tasks/iceberg/delete_rewrite.go b/weed/worker/tasks/iceberg/delete_rewrite.go index 983aeec8a..064dab166 100644 --- a/weed/worker/tasks/iceberg/delete_rewrite.go +++ b/weed/worker/tasks/iceberg/delete_rewrite.go @@ -32,9 +32,14 @@ type deleteRewriteGroup struct { TotalSize int64 } +// positionDeleteRow is the Iceberg position-delete schema: the spec reserves +// field ids 2147483546 (file_path) and 2147483545 (pos). file_path is +// dictionary-encoded: PyArrow cannot decode parquet-go's default +// DELTA_LENGTH_BYTE_ARRAY into the dictionary PyIceberg asks for, and the +// column repeats one data file's path anyway. type positionDeleteRow struct { - FilePath string `parquet:"file_path"` - Pos int64 `parquet:"pos"` + FilePath string `parquet:"file_path,dict,id(2147483546)"` + Pos int64 `parquet:"pos,id(2147483545)"` } func hasEligibleDeleteRewrite( diff --git a/weed/worker/tasks/iceberg/delete_rewrite_schema_test.go b/weed/worker/tasks/iceberg/delete_rewrite_schema_test.go new file mode 100644 index 000000000..58863e835 --- /dev/null +++ b/weed/worker/tasks/iceberg/delete_rewrite_schema_test.go @@ -0,0 +1,54 @@ +package iceberg + +import ( + "bytes" + "testing" + + "github.com/parquet-go/parquet-go" + "github.com/parquet-go/parquet-go/format" +) + +// The position-delete file must carry the spec's field ids and a +// dictionary-encoded file_path: PyIceberg reads delete files with +// file_path as a dictionary column, which PyArrow cannot decode from +// DELTA_LENGTH_BYTE_ARRAY, and by-id readers resolve columns by field id. +func TestPositionDeleteFileHasFieldIDsAndDictionaryPaths(t *testing.T) { + rows := []positionDeleteRow{ + {FilePath: "s3://tb/ns/tbl/data/a.parquet", Pos: 0}, + {FilePath: "s3://tb/ns/tbl/data/a.parquet", Pos: 7}, + {FilePath: "s3://tb/ns/tbl/data/b.parquet", Pos: 3}, + } + data, err := writePositionDeleteFile(rows) + if err != nil { + t.Fatal(err) + } + f, err := parquet.OpenFile(bytes.NewReader(data), int64(len(data))) + if err != nil { + t.Fatal(err) + } + want := map[string]int{"file_path": 2147483546, "pos": 2147483545} + for _, field := range f.Schema().Fields() { + if id, ok := want[field.Name()]; !ok || field.ID() != id { + t.Errorf("field %s id = %d, want %d", field.Name(), field.ID(), id) + } + } + for _, rg := range f.Metadata().RowGroups { + for _, col := range rg.Columns { + if col.MetaData.PathInSchema[0] != "file_path" { + continue + } + dict := false + for _, enc := range col.MetaData.Encoding { + if enc == format.DeltaLengthByteArray { + t.Errorf("file_path is DELTA_LENGTH_BYTE_ARRAY: %v", col.MetaData.Encoding) + } + if enc == format.RLEDictionary || enc == format.PlainDictionary { + dict = true + } + } + if !dict { + t.Errorf("file_path is not dictionary-encoded: %v", col.MetaData.Encoding) + } + } + } +}