mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-10 16:27:47 +02:00
iceberg maintenance: write position-delete files with spec field ids and dictionary paths (#11652)
rewrite_position_delete_files emitted a schema-less parquet file whose file_path column was DELTA_LENGTH_BYTE_ARRAY. PyIceberg reads delete files with file_path as a dictionary column, which PyArrow cannot decode from that encoding, so every table the worker rewrote failed scans outright; readers that resolve columns by field id found none. The row struct now declares the spec's reserved field ids (2147483546 file_path, 2147483545 pos) and dictionary-encodes file_path, which is also the compact choice for a column repeating one data file's path. Generated with [Devin](https://devin.ai) Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
1 parent
4111aa4dd8
commit
f206d021f1
2 files changed
+61
-2
No files matched your search
@@ -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(
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user