Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 35 additions & 17 deletions table/arrow_scanner.go
Original file line number Diff line number Diff line change
Expand Up @@ -530,6 +530,10 @@ func processPositionalDeletes(ctx context.Context, deletes set[int64], cursor *r
}
}

func deletionVectorKeepsRow(keepBits []byte, rowCount, pos int64) bool {
return pos >= rowCount || keepBits[pos>>3]&(1<<(uint(pos)&7)) != 0
}

// filterByDeletionVector returns a pipeline step that drops rows present in
// the bitmap by precomputing a bit-packed keep-mask covering the whole file
// once and slicing the relevant range per batch into compute.FilterRecordBatch.
Expand Down Expand Up @@ -560,16 +564,37 @@ func filterByDeletionVector(ctx context.Context, bitmap *dv.RoaringPositionBitma
currentIdx := nextIdx
nextIdx += nrows

// Wrap (and slice) the shared keep-mask buffer for this batch.
// array.NewSlice on a Boolean array tracks the bit-level offset,
// so we don't need byte-aligned slicing — currentIdx can land
// anywhere within a byte.
full := array.NewBoolean(int(rowCount), buf, nil, 0)
defer full.Release()
sliced := array.NewSlice(full, currentIdx, nextIdx).(*array.Boolean)
defer sliced.Release()
if nextIdx <= rowCount {
// Wrap (and slice) the shared keep-mask buffer for this batch.
// array.NewSlice on a Boolean array tracks the bit-level offset,
// so we don't need byte-aligned slicing — currentIdx can land
// anywhere within a byte.
full := array.NewBoolean(int(rowCount), buf, nil, 0)
defer full.Release()
sliced := array.NewSlice(full, currentIdx, nextIdx).(*array.Boolean)
defer sliced.Release()

return compute.FilterRecordBatch(ctx, r, sliced, compute.DefaultFilterOptions())
}
if currentIdx >= rowCount {
r.Retain()

return r, nil
}

// A stale manifest count can be smaller than the rows emitted by
// the file. Preserve rows beyond the mask rather than slicing past
// its bounds and panicking.
bldr := array.NewBooleanBuilder(mem)
defer bldr.Release()
bldr.Reserve(int(nrows))
for pos := currentIdx; pos < nextIdx; pos++ {
bldr.Append(deletionVectorKeepsRow(keepBits, rowCount, pos))
}
mask := bldr.NewBooleanArray()
defer mask.Release()

return compute.FilterRecordBatch(ctx, r, sliced, compute.DefaultFilterOptions())
return compute.FilterRecordBatch(ctx, r, mask, compute.DefaultFilterOptions())
}

bldr := array.NewBooleanBuilder(mem)
Expand All @@ -580,14 +605,7 @@ func filterByDeletionVector(ctx context.Context, bitmap *dv.RoaringPositionBitma
// keepBits is sized for rowCount; a pos past it means row-group
// metadata and the manifest's File.Count() disagree. Keep the row
// rather than indexing out of bounds.
if pos >= rowCount {
bldr.Append(true)

continue
}
// Test keep bit at absolute position pos: byte pos/8, bit pos%8,
// LSB-first (the layout the fast path's array.NewBoolean reads).
bldr.Append(keepBits[pos>>3]&(1<<(uint(pos)&7)) != 0)
bldr.Append(deletionVectorKeepsRow(keepBits, rowCount, pos))
}
mask := bldr.NewBooleanArray()
defer mask.Release()
Expand Down
36 changes: 36 additions & 0 deletions table/dv_scanner_read_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -347,3 +347,39 @@ func TestFilterByDeletionVectorOutOfBoundsPosition(t *testing.T) {
defer out.Release()
assert.Equal(t, []int64{0, 1, 2}, out.Column(0).(*array.Int64).Int64Values())
}

func TestFilterByDeletionVectorStaleRowCount(t *testing.T) {
ctx := context.Background()
mem := memory.NewGoAllocator()

bitmap := dv.NewRoaringPositionBitmap()
bitmap.Set(1)
bitmap.Set(3)
filter := filterByDeletionVector(ctx, bitmap, 4, (&rowPositionSource{}).cursor())

mkBatch := func(values ...int64) arrow.RecordBatch {
bldr := array.NewInt64Builder(mem)
defer bldr.Release()
bldr.AppendValues(values, nil)
col := bldr.NewArray()
defer col.Release()
schema := arrow.NewSchema([]arrow.Field{{Name: "pos", Type: arrow.PrimitiveTypes.Int64}}, nil)

return array.NewRecordBatch(schema, []arrow.Array{col}, int64(len(values)))
}

withinCount, err := filter(mkBatch(0, 1, 2))
require.NoError(t, err)
defer withinCount.Release()
assert.Equal(t, []int64{0, 2}, withinCount.Column(0).(*array.Int64).Int64Values())

beyondCount, err := filter(mkBatch(3, 4, 5))
require.NoError(t, err)
defer beyondCount.Release()
assert.Equal(t, []int64{4, 5}, beyondCount.Column(0).(*array.Int64).Int64Values())

fullyBeyondCount, err := filter(mkBatch(6, 7))
require.NoError(t, err)
defer fullyBeyondCount.Release()
assert.Equal(t, []int64{6, 7}, fullyBeyondCount.Column(0).(*array.Int64).Int64Values())
}
Loading