zeroshade commented on code in PR #1747:
URL: https://github.com/apache/iceberg-go/pull/1747#discussion_r3833577676
##########
table/inspect_internal_test.go:
##########
@@ -2868,6 +2873,1196 @@ func TestAllEntriesSchemaMatchesEntries(t *testing.T) {
require.True(t,
EntriesSchema(partitionType).Equals(AllEntriesSchema(partitionType)))
}
+func newInspectPositionDeletesMetadata(t *testing.T, formatVersion int)
*MetadataBuilder {
+ return newInspectPositionDeletesMetadataWithSchema(t, formatVersion,
simpleSchema())
+}
+
+func newInspectPositionDeletesMetadataWithSchema(
+ t *testing.T,
+ formatVersion int,
+ tableSchema *iceberg.Schema,
+) *MetadataBuilder {
+ t.Helper()
+
+ mb, err := NewMetadataBuilder(formatVersion)
+ require.NoError(t, err)
+ require.NoError(t, mb.AddSchema(tableSchema))
+ require.NoError(t, mb.SetCurrentSchemaID(0))
+ require.NoError(t, mb.AddPartitionSpec(iceberg.UnpartitionedSpec, true))
+ require.NoError(t, mb.SetDefaultSpecID(0))
+ require.NoError(t, mb.SetLoc("mem://position-deletes/table"))
+ addUnsortedSortOrder(t, mb)
+ metadata, err := mb.Build()
+ require.NoError(t, err)
+ mb, err = MetadataBuilderFromBase(metadata, "")
+ require.NoError(t, err)
+
+ return mb
+}
+
+func inspectPositionDeletesTable(
+ t *testing.T,
+ formatVersion int,
+ mb *MetadataBuilder,
+ fs iceio.WriteFileIO,
+ files []iceberg.DataFile,
+) *Table {
+ return inspectPositionDeletesTableWithSchema(t, formatVersion, mb,
simpleSchema(), fs, files)
+}
+
+func inspectPositionDeletesTableWithSchema(
+ t *testing.T,
+ formatVersion int,
+ mb *MetadataBuilder,
+ tableSchema *iceberg.Schema,
+ fs iceio.WriteFileIO,
+ files []iceberg.DataFile,
+) *Table {
+ t.Helper()
+
+ snapshotID := int64(1)
+ sequenceNumber := int64(1)
+ manifestPath := "mem://position-deletes/table/metadata/deletes.avro"
+ manifestBuffer := &bytes.Buffer{}
+ writer, err := iceberg.NewManifestWriter(
+ formatVersion, manifestBuffer, *iceberg.UnpartitionedSpec,
+ tableSchema, snapshotID,
+
iceberg.WithManifestWriterContent(iceberg.ManifestContentDeletes),
+ )
+ require.NoError(t, err)
+ for _, file := range files {
+ require.NoError(t, writer.Add(iceberg.NewManifestEntry(
+ iceberg.EntryStatusADDED, &snapshotID, &sequenceNumber,
&sequenceNumber, file)))
+ }
+ manifest, err := writer.ToManifestFile(
+ manifestPath, int64(manifestBuffer.Len()),
+ iceberg.WithManifestFileContent(iceberg.ManifestContentDeletes),
+ )
+ require.NoError(t, err)
+ require.NoError(t, fs.WriteFile(manifestPath, manifestBuffer.Bytes()))
+
+ return inspectPositionDeletesTableWithManifests(
+ t, formatVersion, mb, tableSchema, fs,
[]iceberg.ManifestFile{manifest})
+}
+
+func inspectPositionDeletesTableWithManifests(
+ t *testing.T,
+ formatVersion int,
+ mb *MetadataBuilder,
+ tableSchema *iceberg.Schema,
+ fs iceio.WriteFileIO,
+ manifests []iceberg.ManifestFile,
+) *Table {
+ t.Helper()
+
+ snapshotID := int64(1)
+ sequenceNumber := int64(1)
+ manifestListPath := "mem://position-deletes/table/metadata/snap.avro"
+ manifestListBuffer := &bytes.Buffer{}
+ require.NoError(t, iceberg.WriteManifestList(
+ formatVersion, manifestListBuffer, snapshotID, nil,
&sequenceNumber, 0,
+ manifests,
+ ))
+ require.NoError(t, fs.WriteFile(manifestListPath,
manifestListBuffer.Bytes()))
+
+ schemaID := tableSchema.ID
+ snapshot := &Snapshot{
+ SnapshotID: snapshotID,
+ SequenceNumber: sequenceNumber,
+ TimestampMs: time.Now().UnixMilli(),
+ ManifestList: manifestListPath,
+ SchemaID: &schemaID,
+ }
+ if formatVersion >= 3 {
+ firstRowID, addedRows := int64(0), int64(0)
+ snapshot.FirstRowID = &firstRowID
+ snapshot.AddedRows = &addedRows
+ }
+ require.NoError(t, mb.AddSnapshot(snapshot))
+ require.NoError(t, mb.SetSnapshotRef(MainBranch, snapshotID, BranchRef))
+ metadata, err := mb.Build()
+ require.NoError(t, err)
+
+ return New(
+ Identifier{"db", "position_deletes"}, metadata, "metadata.json",
+ func(context.Context) (iceio.IO, error) { return fs, nil }, nil,
+ )
+}
+
+func TestInspectPositionDeletesParquet(t *testing.T) {
+ ctx := context.Background()
+ memFS := iceio.NewMemFS()
+ deletePath := "mem://position-deletes/table/data/delete.parquet"
+ dataPath := "mem://position-deletes/table/data/data.parquet"
+ writePosDeleteParquetToMemFS(t, memFS, deletePath, `[
+ {"file_path": "`+dataPath+`", "pos": 1},
+ {"file_path": "`+dataPath+`", "pos": 3}
+ ]`)
+ deleteFile := newPosDeleteFile(t, deletePath, 2, 128)
+ tbl := inspectPositionDeletesTable(
+ t, 2, newInspectPositionDeletesMetadata(t, 2), memFS,
+ []iceberg.DataFile{deleteFile},
+ )
+
+ rr, err := tbl.Inspect().PositionDeletes(ctx)
+ require.NoError(t, err)
+ defer rr.Release()
+ rec := collectRecord(t, rr)
+ defer rec.Release()
+
+ require.EqualValues(t, 2, rec.NumRows())
+ require.EqualValues(t, 5, rec.NumCols())
+ filePaths := rec.Column(0).(*array.String)
+ require.Equal(t, dataPath, filePaths.Value(0))
+ require.Equal(t, dataPath, filePaths.Value(1))
+ require.Equal(t, []int64{1, 3},
rec.Column(1).(*array.Int64).Int64Values())
+ require.EqualValues(t, 2, rec.Column(2).NullN())
+ require.Equal(t, []int32{0, 0},
rec.Column(3).(*array.Int32).Int32Values())
+ deleteFilePaths := rec.Column(4).(*array.String)
+ require.Equal(t, deletePath, deleteFilePaths.Value(0))
+ require.Equal(t, deletePath, deleteFilePaths.Value(1))
+}
+
+func TestInspectPositionDeletesV3ParquetLeavesDVMetadataNull(t *testing.T) {
+ ctx := context.Background()
+ memFS := iceio.NewMemFS()
+ deletePath := "mem://position-deletes/table/data/delete-v3.parquet"
+ dataPath := "mem://position-deletes/table/data/data.parquet"
+ writePosDeleteParquetToMemFS(t, memFS, deletePath, `[
+ {"file_path": "`+dataPath+`", "pos": 1},
+ {"file_path": "`+dataPath+`", "pos": 3}
+ ]`)
+ const fileSize int64 = 128
+ deleteFileBuilder, err := iceberg.NewDataFileBuilder(
+ *iceberg.UnpartitionedSpec, iceberg.EntryContentPosDeletes,
+ deletePath, iceberg.ParquetFile, nil, nil, nil, 2, fileSize)
+ require.NoError(t, err)
+ deleteFile := deleteFileBuilder.ContentSizeInBytes(fileSize / 2).Build()
+ tbl := inspectPositionDeletesTable(
+ t, 3, newInspectPositionDeletesMetadata(t, 3), memFS,
+ []iceberg.DataFile{deleteFile},
+ )
+
+ rr, err := tbl.Inspect().PositionDeletes(ctx)
+ require.NoError(t, err)
+ defer rr.Release()
+ record := collectRecord(t, rr)
+ defer record.Release()
+
+ require.EqualValues(t, 2, record.NumRows())
+ require.EqualValues(t, 7, record.NumCols())
+ require.EqualValues(t, 2, record.Column(5).NullN())
+ require.EqualValues(t, 2, record.Column(6).NullN())
+}
+
+func TestAppendParquetPositionDeleteRowsRejectsNegativePosition(t *testing.T) {
+ memFS := iceio.NewMemFS()
+ deletePath :=
"mem://position-deletes/table/data/delete-negative.parquet"
+ dataPath := "mem://position-deletes/table/data/data.parquet"
+ writePosDeleteParquetToMemFS(t, memFS, deletePath, `[
+ {"file_path": "`+dataPath+`", "pos": -1}
+ ]`)
+
+ rows := 0
+ keepGoing, err := appendParquetPositionDeleteRows(
+ context.Background(), memFS, newPosDeleteFile(t, deletePath, 1,
128),
+ func(iceberg.DataFile, string, int64, scalar.Scalar, bool)
(bool, error) {
+ rows++
+
+ return true, nil
+ },
+ )
+ require.False(t, keepGoing)
+ require.ErrorContains(t, err, "negative pos -1")
+ require.ErrorIs(t, err, iceberg.ErrInvalidSchema)
+ require.Zero(t, rows)
+}
+
+func TestInspectPositionDeletesDefersLaterManifestReads(t *testing.T) {
+ memFS := iceio.NewMemFS()
+ tableSchema := simpleSchema()
+ deletePath := "mem://position-deletes/table/data/delete-first.parquet"
+ dataPath := "mem://position-deletes/table/data/data.parquet"
+ var deleteRows bytes.Buffer
+ deleteRows.WriteByte('[')
+ for pos := range inspectRecordBatchSize {
+ if pos > 0 {
+ deleteRows.WriteByte(',')
+ }
+ fmt.Fprintf(&deleteRows, `{"file_path": %q, "pos": %d}`,
dataPath, pos)
+ }
+ deleteRows.WriteByte(']')
+ writePosDeleteParquetToMemFS(t, memFS, deletePath, deleteRows.String())
+
+ snapshotID := int64(1)
+ sequenceNumber := int64(1)
+ firstManifest := writeInspectManifest(
+ t, memFS,
"mem://position-deletes/table/metadata/first-deletes.avro",
+ *iceberg.UnpartitionedSpec, tableSchema, snapshotID,
iceberg.ManifestContentDeletes,
+ []iceberg.ManifestEntry{iceberg.NewManifestEntry(
+ iceberg.EntryStatusADDED, &snapshotID, &sequenceNumber,
&sequenceNumber,
+ newPosDeleteFile(t, deletePath, inspectRecordBatchSize,
128),
+ )},
+ )
+ missingManifest := iceberg.NewManifestFile(
+ 2, "mem://position-deletes/table/metadata/not-opened.avro", 1,
0, snapshotID,
+
).Content(iceberg.ManifestContentDeletes).AddedFiles(1).ExistingFiles(0).DeletedFiles(0).Build()
+ tbl := inspectPositionDeletesTableWithManifests(
+ t, 2, newInspectPositionDeletesMetadataWithSchema(t, 2,
tableSchema), tableSchema,
+ memFS, []iceberg.ManifestFile{firstManifest, missingManifest},
+ )
+
+ rr, err := tbl.Inspect().PositionDeletes(context.Background())
+ require.NoError(t, err)
+ defer rr.Release()
+ require.True(t, rr.Next(), "reader error: %v", rr.Err())
+ require.EqualValues(t, inspectRecordBatchSize,
rr.RecordBatch().NumRows())
+}
+
+func TestAppendParquetPositionDeleteRowsRejectsNullRow(t *testing.T) {
+ memFS := iceio.NewMemFS()
+ deletePath :=
"mem://position-deletes/table/data/delete-null-row.parquet"
+ dataPath := "mem://position-deletes/table/data/data.parquet"
+ schema := arrow.NewSchema([]arrow.Field{
+ {Name: "file_path", Type: arrow.BinaryTypes.String, Nullable:
false},
+ {Name: "pos", Type: arrow.PrimitiveTypes.Int64, Nullable:
false},
+ {Name: "row", Type: arrow.StructOf(arrow.Field{
+ Name: "id",
+ Type: arrow.PrimitiveTypes.Int32,
+ Nullable: false,
+ }), Nullable: true},
+ }, nil)
+ bldr := array.NewRecordBuilder(memory.DefaultAllocator, schema)
+ bldr.Field(0).(*array.StringBuilder).Append(dataPath)
+ bldr.Field(1).(*array.Int64Builder).Append(1)
+ bldr.Field(2).(*array.StructBuilder).Append(false)
+ record := bldr.NewRecordBatch()
+ defer record.Release()
+ defer bldr.Release()
+
+ tbl := array.NewTableFromRecords(schema, []arrow.RecordBatch{record})
+ defer tbl.Release()
+ file, err := memFS.Create(deletePath)
+ require.NoError(t, err)
+ require.NoError(t, pqarrow.WriteTable(
+ tbl, file, record.NumRows(),
+ parquet.NewWriterProperties(parquet.WithStats(true)),
pqarrow.DefaultWriterProps()))
+ require.NoError(t, file.Close())
+
+ rows := 0
+ keepGoing, err := appendParquetPositionDeleteRows(
+ context.Background(), memFS, newPosDeleteFile(t, deletePath, 1,
128),
+ func(iceberg.DataFile, string, int64, scalar.Scalar, bool)
(bool, error) {
+ rows++
+
+ return true, nil
+ },
+ )
+ require.False(t, keepGoing)
+ require.ErrorIs(t, err, iceberg.ErrInvalidSchema)
+ require.ErrorContains(t, err, "null row")
+ require.Zero(t, rows)
+}
+
+func TestInspectPositionDeletesParquetProjectsEvolvedNestedRow(t *testing.T) {
Review Comment:
Please make this nested evolution test use
`WithInspectAllocator(memory.NewCheckedAllocator(...))` and assert zero bytes
remain. The existing value assertions pass while the scan leaks 320 bytes, so
the checked allocator is needed to prevent this ownership regression.
##########
table/inspect_position_deletes.go:
##########
@@ -0,0 +1,865 @@
+// 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"
+ "errors"
+ "fmt"
+ "math"
+
+ "github.com/apache/arrow-go/v18/arrow"
+ "github.com/apache/arrow-go/v18/arrow/array"
+ "github.com/apache/arrow-go/v18/arrow/compute"
+ "github.com/apache/arrow-go/v18/arrow/scalar"
+ "github.com/apache/iceberg-go"
+ iceio "github.com/apache/iceberg-go/io"
+ "github.com/apache/iceberg-go/table/dv"
+ tblutils "github.com/apache/iceberg-go/table/internal"
+)
+
+const (
+ positionDeleteFilePathID = math.MaxInt32 - 101
+ positionDeletePosID = math.MaxInt32 - 102
+ positionDeleteRowID = math.MaxInt32 - 103
+ positionDeletePartitionID = math.MaxInt32 - 5
+ positionDeleteSpecID = math.MaxInt32 - 4
+ positionDeletePhysicalPathID = math.MaxInt32 - 1
+ positionDeleteContentOffsetID = math.MaxInt32 - 6
+ positionDeleteContentSizeID = math.MaxInt32 - 7
+ positionDeletePhysicalPathName = "delete_file_path"
+)
+
+// PositionDeletes returns the individual position-delete records referenced
+// by the current snapshot. Parquet position-delete files and V3 deletion
+// vectors are exposed through the same schema.
+func (i InspectTable) PositionDeletes(ctx context.Context)
(array.RecordReader, error) {
+ partitionType, partitionIDs, err :=
positionDeletesPartitionType(i.tbl.metadata)
+ if err != nil {
+ return nil, fmt.Errorf("inspect position deletes: %w", err)
+ }
+ schema := PositionDeletesSchema(i.tbl.metadata.CurrentSchema(),
partitionType, i.tbl.metadata.Version())
+ arrowSchema, err := SchemaToArrowSchema(schema, nil, true, false)
+ if err != nil {
+ return nil, fmt.Errorf("inspect position deletes: build arrow
schema: %w", err)
+ }
+
+ fs, manifests, err := i.currentPositionDeleteManifests(ctx)
+ if err != nil {
+ return nil, fmt.Errorf("inspect position deletes: %w", err)
+ }
+ ctx = compute.WithAllocator(ctx, i.alloc)
+
+ return i.positionDeleteRecordReader(
+ ctx, arrowSchema, fs, manifests, partitionType, partitionIDs,
i.tbl.metadata.Version()), nil
+}
+
+func (i InspectTable) currentPositionDeleteManifests(
+ ctx context.Context,
+) (iceio.IO, []iceberg.ManifestFile, error) {
+ snapshot := i.tbl.metadata.CurrentSnapshot()
+ if snapshot == nil {
+ return nil, nil, nil
+ }
+ if i.tbl.fsF == nil {
+ return nil, nil, errors.New("table file IO is not configured")
+ }
+ fs, err := i.tbl.fsF(ctx)
+ if err != nil {
+ return nil, nil, err
+ }
+ manifests, err := snapshot.Manifests(fs)
+ if err != nil {
+ return nil, nil, err
+ }
+
+ return fs, manifests, nil
+}
+
+type positionDeleteRecordAppender struct {
+ filePath *array.StringBuilder
+ pos *array.Int64Builder
+ row *array.StructBuilder
+ partition *inspectPartitionBuilder
+ specID *array.Int32Builder
+ deleteFilePath *array.StringBuilder
+ contentOffset *array.Int64Builder
+ contentSize *array.Int64Builder
+ partitionType *iceberg.StructType
+ partitionIDByOld map[int]int
+ formatVersion int
+ projection positionDeleteProjectionOptions
+}
+
+type positionDeleteProjectionOptions struct {
+ ctx context.Context
+ formatVersion int
+ tableSchema *iceberg.Schema
+ nameMapping iceberg.NameMapping
+ tableProperties iceberg.Properties
+}
+
+func newPositionDeleteRecordAppender(
+ bldr *array.RecordBuilder,
+ partitionType *iceberg.StructType,
+ partitionIDByOld map[int]int,
+ formatVersion int,
+ options ...positionDeleteProjectionOptions,
+) (positionDeleteRecordAppender, error) {
+ nextField := 3
+ projection := positionDeleteProjectionOptions{
+ ctx: context.Background(),
+ formatVersion: formatVersion,
+ }
+ if len(options) > 1 {
+ return positionDeleteRecordAppender{}, fmt.Errorf("%w: expected
at most one position delete projection option", iceberg.ErrInvalidArgument)
+ }
+ if len(options) == 1 {
+ projection = options[0]
+ if projection.ctx == nil {
+ projection.ctx = context.Background()
+ }
+ if projection.formatVersion == 0 {
+ projection.formatVersion = formatVersion
+ }
+ }
+ if projection.tableSchema == nil {
+ rowType, ok := bldr.Field(2).Type().(*arrow.StructType)
+ if !ok {
+ return positionDeleteRecordAppender{}, fmt.Errorf("%w:
destination row has type %s, want struct", iceberg.ErrInvalidSchema,
bldr.Field(2).Type())
+ }
+ rowSchema, err := ArrowSchemaToIcebergWithOptions(
+ arrow.NewSchema(rowType.Fields(), nil),
+ ArrowToIcebergOptions{TableProperties:
projection.tableProperties},
+ )
+ if err != nil {
+ return positionDeleteRecordAppender{}, fmt.Errorf("%w:
resolve destination row schema: %v", iceberg.ErrInvalidSchema, err)
+ }
+ projection.tableSchema = rowSchema
+ }
+ if projection.nameMapping == nil && projection.tableSchema != nil {
+ projection.nameMapping = projection.tableSchema.NameMapping()
+ }
+ out := positionDeleteRecordAppender{
+ filePath: bldr.Field(0).(*array.StringBuilder),
+ pos: bldr.Field(1).(*array.Int64Builder),
+ row: bldr.Field(2).(*array.StructBuilder),
+ partitionType: partitionType,
+ partitionIDByOld: partitionIDByOld,
+ formatVersion: formatVersion,
+ projection: projection,
+ }
+ if len(partitionType.FieldList) > 0 {
+ partitionBuilder, err := newInspectPartitionBuilder(
+ bldr.Field(nextField).(*array.StructBuilder),
partitionType)
+ if err != nil {
+ return out, err
+ }
+ out.partition = partitionBuilder
+ nextField++
+ }
+ out.specID = bldr.Field(nextField).(*array.Int32Builder)
+ out.deleteFilePath = bldr.Field(nextField + 1).(*array.StringBuilder)
+ if formatVersion >= 3 {
+ out.contentOffset = bldr.Field(nextField +
2).(*array.Int64Builder)
+ out.contentSize = bldr.Field(nextField +
3).(*array.Int64Builder)
+ }
+
+ return out, nil
+}
+
+func (a positionDeleteRecordAppender) append(
+ file iceberg.DataFile,
+ dataFilePath string,
+ pos int64,
+ deletedRow scalar.Scalar,
+ projected ...bool,
+) error {
+ a.filePath.Append(dataFilePath)
+ a.pos.Append(pos)
+ var err error
+ if len(projected) > 0 && projected[0] {
+ if deletedRow == nil || !deletedRow.IsValid() {
+ a.row.AppendNull()
+ } else {
+ err = appendPositionDeleteProjectedScalar(a.row,
deletedRow)
+ }
+ } else {
+ err = appendProjectedPositionDeleteRow(a.row, deletedRow,
a.projection)
+ }
+ if err != nil {
+ return fmt.Errorf("append deleted row: %w", err)
+ }
+
+ if a.partition != nil {
+ partition := make(map[int]any, len(file.Partition()))
+ for oldID, value := range file.Partition() {
+ if newID, ok := a.partitionIDByOld[oldID]; ok {
+ partition[newID] = value
+ }
+ }
+ if err := a.partition.append(partition); err != nil {
+ return err
+ }
+ }
+ a.specID.Append(file.SpecID())
+ a.deleteFilePath.Append(file.FilePath())
+ if a.formatVersion >= 3 {
+ appendInspectOptionalInt64(a.contentOffset,
file.ContentOffset())
+ appendInspectOptionalInt64(a.contentSize,
positionDeleteContentSize(file))
+ }
+
+ return nil
+}
+
+func positionDeleteContentSize(file iceberg.DataFile) *int64 {
+ if file.FileFormat() != iceberg.PuffinFile {
+ return nil
+ }
+
+ return file.ContentSizeInBytes()
+}
+
+// appendPositionDeleteRow projects a row from a position-delete file onto the
+// current table schema. The shared Arrow projection path resolves field IDs,
+// name mappings, nested types, promotions, and initial defaults consistently
+// with normal table scans.
+func appendProjectedPositionDeleteRow(
+ builder *array.StructBuilder,
+ deletedRow scalar.Scalar,
+ projection positionDeleteProjectionOptions,
+) error {
+ if deletedRow == nil || !deletedRow.IsValid() {
+ builder.AppendNull()
+
+ return nil
+ }
+
+ row, ok := deletedRow.(*scalar.Struct)
+ if !ok {
+ return fmt.Errorf("%w: row has type %s, want struct",
iceberg.ErrInvalidSchema, deletedRow.DataType())
+ }
+ rowType, ok := row.DataType().(*arrow.StructType)
+ if !ok {
+ return fmt.Errorf("%w: row has type %s, want struct",
iceberg.ErrInvalidSchema, row.DataType())
+ }
+ rowBuilder :=
array.NewStructBuilder(compute.GetAllocator(projection.ctx), rowType)
+ defer rowBuilder.Release()
+ if err := appendPositionDeleteProjectedScalar(rowBuilder, row); err !=
nil {
+ return fmt.Errorf("make row array: %w", err)
+ }
+ rowArray := rowBuilder.NewArray()
+ defer rowArray.Release()
+
+ rowStruct, ok := rowArray.(*array.Struct)
+ if !ok {
+ return fmt.Errorf("%w: row array has type %s, want struct",
iceberg.ErrInvalidSchema, rowArray.DataType())
+ }
+ projected, err := projectPositionDeleteRows(projection.ctx, rowStruct,
projection)
+ if err != nil {
+ return err
+ }
+ defer projected.Release()
+
+ projectedStruct := array.RecordToStructArray(projected)
+ defer projectedStruct.Release()
+ projectedRow, err := scalar.GetScalar(projectedStruct, 0)
+ if err != nil {
+ return err
+ }
+ if releasable, ok := projectedRow.(scalar.Releasable); ok {
+ defer releasable.Release()
+ }
+
+ return appendPositionDeleteProjectedScalar(builder, projectedRow)
+}
+
+func appendPositionDeleteProjectedScalar(builder array.Builder, value
scalar.Scalar) error {
+ if value == nil || !value.IsValid() {
+ builder.AppendNull()
+
+ return nil
+ }
+
+ switch builder := builder.(type) {
+ case *array.StructBuilder:
+ source, ok := value.(*scalar.Struct)
+ if !ok {
+ return fmt.Errorf("%w: projected struct has type %s,
want struct",
+ iceberg.ErrInvalidSchema, value.DataType())
+ }
+ if len(source.Value) != builder.NumField() {
+ return fmt.Errorf("%w: projected struct has %d fields,
want %d",
+ iceberg.ErrInvalidSchema, len(source.Value),
builder.NumField())
+ }
+ builder.Append(true)
+ for index, child := range source.Value {
+ if err :=
appendPositionDeleteProjectedScalar(builder.FieldBuilder(index), child); err !=
nil {
+ return fmt.Errorf("projected struct field %d:
%w", index, err)
+ }
+ }
+
+ return nil
+ case *array.MapBuilder:
+ source, ok := value.(*scalar.Map)
+ if !ok {
+ return fmt.Errorf("%w: projected map has type %s",
iceberg.ErrInvalidSchema, value.DataType())
+ }
+ entries := source.GetList()
+ if entries == nil {
+ return fmt.Errorf("%w: projected map has no entries",
iceberg.ErrInvalidSchema)
+ }
+ entryType, ok := entries.DataType().(*arrow.StructType)
+ if !ok || entryType.NumFields() != 2 {
+ return fmt.Errorf("%w: projected map entries have type
%s", iceberg.ErrInvalidSchema, entries.DataType())
+ }
+ builder.Append(true)
+ for index := range entries.Len() {
+ entry, err := scalar.GetScalar(entries, index)
+ if err != nil {
+ return err
+ }
+ entryStruct, ok := entry.(*scalar.Struct)
+ if !ok || len(entryStruct.Value) != 2 {
+ if releasable, ok := entry.(scalar.Releasable);
ok {
+ releasable.Release()
+ }
+
+ return fmt.Errorf("%w: projected map entry %d
is not a two-field struct",
+ iceberg.ErrInvalidSchema, index)
+ }
+ if entryStruct.Value[0] == nil ||
!entryStruct.Value[0].IsValid() {
+ if releasable, ok := entry.(scalar.Releasable);
ok {
+ releasable.Release()
+ }
+
+ return fmt.Errorf("%w: projected map entry %d
has a null key",
+ iceberg.ErrInvalidSchema, index)
+ }
+ if err :=
appendPositionDeleteProjectedScalar(builder.KeyBuilder(),
entryStruct.Value[0]); err != nil {
+ if releasable, ok := entry.(scalar.Releasable);
ok {
+ releasable.Release()
+ }
+
+ return fmt.Errorf("map entry %d key: %w",
index, err)
+ }
+ if err :=
appendPositionDeleteProjectedScalar(builder.ItemBuilder(),
entryStruct.Value[1]); err != nil {
+ if releasable, ok := entry.(scalar.Releasable);
ok {
+ releasable.Release()
+ }
+
+ return fmt.Errorf("map entry %d value: %w",
index, err)
+ }
+ if releasable, ok := entry.(scalar.Releasable); ok {
+ releasable.Release()
+ }
+ }
+
+ return nil
+ case array.ListLikeBuilder:
+ source, ok := value.(scalar.ListScalar)
+ if !ok {
+ return fmt.Errorf("%w: projected list has type %s",
iceberg.ErrInvalidSchema, value.DataType())
+ }
+ values := source.GetList()
+ if values == nil {
+ return fmt.Errorf("%w: projected list has no values",
iceberg.ErrInvalidSchema)
+ }
+ if destination, ok :=
builder.Type().(*arrow.FixedSizeListType); ok &&
+ values.Len() != int(destination.Len()) {
+ return fmt.Errorf("%w: projected list has %d elements,
want %d",
+ iceberg.ErrInvalidSchema, values.Len(),
destination.Len())
+ }
+ builder.Append(true)
+ for index := range values.Len() {
+ child, err := scalar.GetScalar(values, index)
+ if err != nil {
+ return err
+ }
+ appendErr :=
appendPositionDeleteProjectedScalar(builder.ValueBuilder(), child)
+ if releasable, ok := child.(scalar.Releasable); ok {
+ releasable.Release()
+ }
+ if appendErr != nil {
+ return fmt.Errorf("list element %d: %w", index,
appendErr)
+ }
+ }
+
+ return nil
+ default:
+ return scalar.Append(builder, value)
+ }
+}
+
+func projectPositionDeleteRows(
+ ctx context.Context,
+ rows *array.Struct,
+ projection positionDeleteProjectionOptions,
+) (arrow.RecordBatch, error) {
+ if projection.tableSchema == nil {
+ return nil, fmt.Errorf("%w: position delete projection has no
table schema", iceberg.ErrInvalidSchema)
+ }
+ if rows == nil {
+ return nil, fmt.Errorf("%w: position delete row array is nil",
iceberg.ErrInvalidSchema)
+ }
+ rowType, ok := rows.DataType().(*arrow.StructType)
+ if !ok {
+ return nil, fmt.Errorf("%w: row has type %s, want struct",
iceberg.ErrInvalidSchema, rows.DataType())
+ }
+
+ sourceArrowSchema := arrow.NewSchema(rowType.Fields(), nil)
+ fileSchema, err := ArrowSchemaToIcebergWithOptions(sourceArrowSchema,
ArrowToIcebergOptions{
+ NameMapping: projection.nameMapping,
+ TableSchema: projection.tableSchema,
+ TableProperties: projection.tableProperties,
+ })
+ if err != nil {
+ return nil, fmt.Errorf("resolve position delete row schema:
%w", err)
+ }
+
+ sourceBatch := array.RecordFromStructArray(rows, sourceArrowSchema)
+ defer sourceBatch.Release()
+
+ return ToRequestedSchema(ctx, projection.tableSchema, fileSchema,
sourceBatch, SchemaOptions{
+ FormatVersion: projection.formatVersion,
+ IncludeFieldIDs: true,
+ AllowMissingRequired: true,
+ TableProperties: projection.tableProperties,
+ })
+}
+
+func (i InspectTable) positionDeleteRecordReader(
+ ctx context.Context,
+ arrowSchema *arrow.Schema,
+ fs iceio.IO,
+ manifests []iceberg.ManifestFile,
+ partitionType *iceberg.StructType,
+ partitionIDByOld map[int]int,
+ formatVersion int,
+) array.RecordReader {
+ return array.ReaderFromIter(arrowSchema, func(yield
func(arrow.RecordBatch, error) bool) {
+ if err := ctx.Err(); err != nil {
+ _ = yield(nil, err)
+
+ return
+ }
+
+ bldr := array.NewRecordBuilder(i.alloc, arrowSchema)
+ defer bldr.Release()
+ appender, err := newPositionDeleteRecordAppender(
+ bldr, partitionType, partitionIDByOld, formatVersion,
+ positionDeleteProjectionOptions{
+ ctx: ctx,
+ formatVersion: formatVersion,
+ tableSchema: i.tbl.metadata.CurrentSchema(),
+ nameMapping: i.tbl.metadata.NameMapping(),
+ tableProperties: i.tbl.metadata.Properties(),
+ },
+ )
+ if err != nil {
+ _ = yield(nil, err)
+
+ return
+ }
+
+ rows := 0
+ emitted := false
+ emit := func() bool {
+ if rows == 0 {
+ return true
+ }
+ batch := bldr.NewRecordBatch()
+ rows = 0
+ emitted = true
+
+ return yield(batch, nil)
+ }
+ yieldError := func(err error) {
+ _ = yield(nil, err)
+ }
+ appendRow := func(file iceberg.DataFile, path string, pos
int64, deletedRow scalar.Scalar, projected bool) (bool, error) {
+ if err := appender.append(file, path, pos, deletedRow,
projected); err != nil {
+ return false, err
+ }
+ rows++
+ if rows == inspectRecordBatchSize {
+ return emit(), nil
+ }
+
+ return true, nil
+ }
+
+ for _, manifest := range manifests {
+ if err := ctx.Err(); err != nil {
+ yieldError(err)
+
+ return
+ }
+ if manifest.ManifestContent() !=
iceberg.ManifestContentDeletes {
+ continue
+ }
+
+ for entry, err := range manifest.Entries(fs, true) {
+ if err != nil {
+ yieldError(fmt.Errorf("read manifest
%s: %w", manifest.FilePath(), err))
+
+ return
+ }
+ if err := ctx.Err(); err != nil {
+ yieldError(err)
+
+ return
+ }
+ file := entry.DataFile()
+ if file.ContentType() !=
iceberg.EntryContentPosDeletes {
+ continue
+ }
+
+ var keepGoing bool
+ switch file.FileFormat() {
+ case iceberg.PuffinFile:
+ keepGoing, err =
appendDeletionVectorRows(ctx, fs, file, appendRow)
+ case iceberg.ParquetFile:
+ keepGoing, err =
appendParquetPositionDeleteRows(ctx, fs, file, appendRow, appender.projection)
+ default:
+ keepGoing = false
+ err = fmt.Errorf("%w: unsupported
position delete file format %s",
+ iceberg.ErrNotImplemented,
file.FileFormat())
+ }
+ if err != nil {
+ yieldError(fmt.Errorf("read position
delete file %s: %w", file.FilePath(), err))
+
+ return
+ }
+ if !keepGoing {
+ return
+ }
+ }
+ }
+
+ if rows > 0 {
+ _ = emit()
+ } else if !emitted {
+ batch := bldr.NewRecordBatch()
+ _ = yield(batch, nil)
+ }
+ })
+}
+
+type appendPositionDeleteRow func(iceberg.DataFile, string, int64,
scalar.Scalar, bool) (bool, error)
+
+func appendDeletionVectorRows(
+ ctx context.Context,
+ fs iceio.IO,
+ file iceberg.DataFile,
+ appendRow appendPositionDeleteRow,
+) (bool, error) {
+ referencedDataFile := file.ReferencedDataFile()
+ if referencedDataFile == nil {
+ return false, fmt.Errorf("%w: deletion vector is missing
referenced_data_file",
+ iceberg.ErrInvalidSchema)
+ }
+ if file.ContentOffset() == nil || file.ContentSizeInBytes() == nil {
+ return false, fmt.Errorf("%w: deletion vector is missing
content_offset/content_size_in_bytes",
+ iceberg.ErrInvalidSchema)
+ }
+ bitmap, err := dv.ReadDV(fs, file)
+ if err != nil {
+ return false, err
+ }
+ for position := range bitmap.Positions() {
+ if err := ctx.Err(); err != nil {
+ return false, err
+ }
+ if position > math.MaxInt64 {
+ return false, fmt.Errorf("%w: deletion position %d
exceeds int64", iceberg.ErrInvalidSchema, position)
+ }
+ keepGoing, err := appendRow(file, *referencedDataFile,
int64(position), nil, false)
+ if err != nil || !keepGoing {
+ return keepGoing, err
+ }
+ }
+
+ return true, nil
+}
+
+func appendParquetPositionDeleteRows(
+ ctx context.Context,
+ fs iceio.IO,
+ file iceberg.DataFile,
+ appendRow appendPositionDeleteRow,
+ options ...positionDeleteProjectionOptions,
+) (keepGoing bool, err error) {
+ projection := positionDeleteProjectionOptions{ctx: ctx}
+ if len(options) > 1 {
+ return false, fmt.Errorf("%w: expected at most one position
delete projection option", iceberg.ErrInvalidArgument)
+ }
+ if len(options) == 1 {
+ projection = options[0]
+ if projection.ctx == nil {
+ projection.ctx = ctx
+ }
+ }
+ source, err := tblutils.GetFile(ctx, fs, file, true)
+ if err != nil {
+ return false, err
+ }
+ reader, err := source.GetReader(ctx)
+ if err != nil {
+ return false, err
+ }
+ defer func() {
+ if closeErr := reader.Close(); err == nil && closeErr != nil {
+ err = closeErr
+ }
+ }()
+
+ records, err := reader.GetRecords(ctx, nil, nil)
+ if err != nil {
+ return false, err
+ }
+ defer records.Release()
+
+ for records.Next() {
+ continueReading, appendErr := appendPositionDeleteRecord(
+ ctx, file, records.RecordBatch(), appendRow, projection)
+ if appendErr != nil || !continueReading {
+ return continueReading, appendErr
+ }
+ }
+ if err := records.Err(); err != nil {
+ return false, err
+ }
+
+ return true, nil
+}
+
+func appendPositionDeleteRecord(
+ ctx context.Context,
+ file iceberg.DataFile,
+ record arrow.RecordBatch,
+ appendRow appendPositionDeleteRow,
+ projection positionDeleteProjectionOptions,
+) (bool, error) {
+ filePathIndex, posIndex, err :=
positionDeleteColumnIndices(record.Schema())
+ if err != nil {
+ return false, err
+ }
+ filePaths, err := filePathValues(record.Column(filePathIndex))
+ if err != nil {
+ return false, err
+ }
+ if err := validatePositionDeleteFilePathValues(filePaths,
record.Column(filePathIndex)); err != nil {
+ return false, err
+ }
+ positions, ok := record.Column(posIndex).(*array.Int64)
+ if !ok {
+ return false, fmt.Errorf("%w: pos column has type %s, want
int64",
+ iceberg.ErrInvalidSchema,
record.Column(posIndex).DataType())
+ }
+ if positions.NullN() > 0 {
+ return false, fmt.Errorf("%w: null pos in position delete
file", iceberg.ErrInvalidSchema)
+ }
+ if err := validatePositionDeletePositions(positions); err != nil {
+ return false, err
+ }
+
+ rowIndex := -1
+ if indices := record.Schema().FieldIndices("row"); len(indices) > 1 {
+ return false, fmt.Errorf("%w: position delete file contains
multiple row columns",
+ iceberg.ErrInvalidSchema)
+ } else if len(indices) == 1 {
+ rowIndex = indices[0]
+ }
+ if rowIndex >= 0 && record.Column(rowIndex).NullN() > 0 {
+ return false, fmt.Errorf("%w: null row in position delete
file", iceberg.ErrInvalidSchema)
+ }
+
+ var projectedRecord arrow.RecordBatch
+ var projectedRows *array.Struct
+ if rowIndex >= 0 && projection.tableSchema != nil {
+ rowArray, ok := record.Column(rowIndex).(*array.Struct)
+ if !ok {
+ return false, fmt.Errorf("%w: row column has type %s,
want struct",
+ iceberg.ErrInvalidSchema,
record.Column(rowIndex).DataType())
+ }
+ projectedRecord, err = projectPositionDeleteRows(ctx, rowArray,
projection)
+ if err != nil {
+ return false, err
+ }
+ projectedRows = array.RecordToStructArray(projectedRecord)
+ defer projectedRecord.Release()
+ defer projectedRows.Release()
+ }
+
+ for row := range int(record.NumRows()) {
+ if err := ctx.Err(); err != nil {
+ return false, err
+ }
+ var deletedRow scalar.Scalar
+ projected := false
+ if rowIndex >= 0 {
+ if projectedRows != nil {
+ deletedRow, err =
scalar.GetScalar(projectedRows, row)
Review Comment:
Blocking: the nested projected-row path still leaks Arrow allocations. I
reran `TestInspectPositionDeletesParquetProjectsEvolvedNestedRow` with only
`Inspect` wired to a checked allocator; after releasing the reader and returned
record, the one-row scan retained 320 bytes (three builder allocations and two
compute-kernel allocations, each 64 bytes). Please balance the ownership in
this projection/scalar lifecycle so the checked allocator returns to zero.
--
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]