Skip to content
Draft
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
9 changes: 5 additions & 4 deletions src/iceberg/deletes/dv_writer.cc
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,6 @@

#include "iceberg/deletes/dv_util_internal.h"
#include "iceberg/deletes/position_delete_index.h"
#include "iceberg/deletes/roaring_position_bitmap.h"
#include "iceberg/file_format.h"
#include "iceberg/file_io.h" // IWYU pragma: keep
#include "iceberg/manifest/manifest_entry.h"
Expand Down Expand Up @@ -79,9 +78,7 @@ class DVWriter::Impl {
ICEBERG_PRECHECK(!referenced_data_file.empty(),
"Deletion vector requires a non-empty referenced data file");
ICEBERG_PRECHECK(spec != nullptr, "Deletion vector requires a partition spec");
ICEBERG_PRECHECK(pos >= 0 && pos <= RoaringPositionBitmap::kMaxPosition,
"Deletion vector position out of range [0, {}]: {}",
RoaringPositionBitmap::kMaxPosition, pos);
ICEBERG_PRECHECK(pos >= 0, "Deletion vector position must be non-negative: {}", pos);
DeletesFor(referenced_data_file, spec, partition).positions.Delete(pos);
return {};
}
Expand Down Expand Up @@ -112,6 +109,10 @@ class DVWriter::Impl {
ICEBERG_RETURN_UNEXPECTED(LoadPreviousDeletes(path, deletes));
}

for (auto& [_, deletes] : deletes_by_path_) {
ICEBERG_RETURN_UNEXPECTED(deletes.positions.PrepareForSerialization());
}

ICEBERG_ASSIGN_OR_RAISE(auto output_file, options_.io->NewOutputFile(options_.path));
const std::string output_path(options_.path);
ICEBERG_ASSIGN_OR_RAISE(
Expand Down
19 changes: 13 additions & 6 deletions src/iceberg/deletes/position_delete_index.cc
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ constexpr std::array<uint8_t, 4> kMagic = {0xD1, 0xD3, 0x39, 0x64};
constexpr int32_t kLengthPrefixBytes = 4;
constexpr int32_t kMagicBytes = 4;
constexpr int32_t kCrcBytes = 4;
constexpr size_t kMaxSerializedLength = std::numeric_limits<int32_t>::max();

uint32_t ComputeCrc32(std::span<const uint8_t> bytes) {
uLong crc = crc32(0L, Z_NULL, 0);
Expand Down Expand Up @@ -142,16 +143,22 @@ void PositionDeleteIndex::Merge(const PositionDeleteIndex& other) {
other.delete_files_.end());
}

Result<std::vector<uint8_t>> PositionDeleteIndex::Serialize() {
Result<size_t> PositionDeleteIndex::PrepareForSerialization() {
bitmap_.Optimize(); // run-length encode before serializing
std::vector<uint8_t> blob(kLengthPrefixBytes);
blob.insert(blob.end(), kMagic.begin(), kMagic.end());
ICEBERG_ASSIGN_OR_RAISE(const auto vector_size, bitmap_.SerializeTo(blob));

const size_t vector_size = bitmap_.SerializedSizeInBytes();
const size_t magic_and_vector_size = kMagicBytes + vector_size;
ICEBERG_PRECHECK(magic_and_vector_size <= std::numeric_limits<int32_t>::max(),
ICEBERG_PRECHECK(magic_and_vector_size <= kMaxSerializedLength,
"Deletion vector is too large to serialize: {} bytes",
magic_and_vector_size);
return magic_and_vector_size;
}

Result<std::vector<uint8_t>> PositionDeleteIndex::Serialize() {
ICEBERG_ASSIGN_OR_RAISE(const auto magic_and_vector_size, PrepareForSerialization());

std::vector<uint8_t> blob(kLengthPrefixBytes);
blob.insert(blob.end(), kMagic.begin(), kMagic.end());
ICEBERG_RETURN_UNEXPECTED(bitmap_.SerializeTo(blob));

WriteBigEndian(static_cast<int32_t>(magic_and_vector_size), blob.data());
const auto crc_offset = blob.size();
Expand Down
8 changes: 8 additions & 0 deletions src/iceberg/deletes/position_delete_index.h
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
/// \file iceberg/deletes/position_delete_index.h
/// Index of deleted row positions for a data file.

#include <cstddef>
#include <cstdint>
#include <memory>
#include <span>
Expand All @@ -34,6 +35,8 @@

namespace iceberg {

class DVWriter;

/// \brief Tracks deleted row positions using a bitmap.
///
/// This class provides a domain-specific API for position deletes
Expand All @@ -53,6 +56,8 @@ class ICEBERG_EXPORT PositionDeleteIndex {
/// \brief Mark a range of positions as deleted [pos_start, pos_end).
/// \param pos_start Start position (inclusive)
/// \param pos_end End position (exclusive)
/// \note Because pos_end is an int64_t exclusive endpoint, this method cannot
/// include INT64_MAX. Call Delete(INT64_MAX) separately.
void Delete(int64_t pos_start, int64_t pos_end);

/// \brief Check if a position is deleted.
Expand Down Expand Up @@ -97,6 +102,8 @@ class ICEBERG_EXPORT PositionDeleteIndex {
private:
explicit PositionDeleteIndex(RoaringPositionBitmap bitmap);

Result<size_t> PrepareForSerialization();

// Bulk-add positions sharing high-32-bit `key`. Private hook for
// `ForEachPositionDelete`'s bulk path; keeps `Delete` the sole public
// mutation surface.
Expand All @@ -105,6 +112,7 @@ class ICEBERG_EXPORT PositionDeleteIndex {
friend void ICEBERG_EXPORT ForEachPositionDelete(std::span<const int64_t> positions,
PositionDeleteIndex& target,
std::vector<uint32_t>& scratch);
friend class DVWriter;

RoaringPositionBitmap bitmap_;
std::vector<std::shared_ptr<DataFile>> delete_files_;
Expand Down
10 changes: 5 additions & 5 deletions src/iceberg/deletes/position_delete_range_consumer.cc
Original file line number Diff line number Diff line change
Expand Up @@ -31,9 +31,7 @@ namespace iceberg {

namespace {

bool IsValidPosition(int64_t pos) {
return pos >= 0 && pos <= RoaringPositionBitmap::kMaxPosition;
}
bool IsValidPosition(int64_t pos) { return pos >= 0; }

// Unsigned subtraction so negative or wrap-around input can't
// false-positive via signed overflow.
Expand All @@ -45,11 +43,13 @@ bool IsAdjacent(int64_t prev, int64_t next) {
// bulk path groups by this key before flushing via `BulkAddForKey`.
int32_t HighKeyFromPosition(int64_t pos) { return static_cast<int32_t>(pos >> 32); }

// Emit `[range_start, last_position]`, collapsing singletons. Callers
// pre-filter via `IsValidPosition`, so `last_position + 1` cannot overflow.
// Emit `[range_start, last_position]`, collapsing singletons.
void EmitRange(PositionDeleteIndex& target, int64_t range_start, int64_t last_position) {
if (range_start == last_position) {
target.Delete(range_start);
} else if (last_position == RoaringPositionBitmap::kMaxPosition) {
target.Delete(range_start, last_position);
target.Delete(last_position);
} else {
target.Delete(range_start, last_position + 1);
}
Expand Down
73 changes: 27 additions & 46 deletions src/iceberg/deletes/roaring_position_bitmap.cc
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
#include <cstring>
#include <exception>
#include <limits>
#include <map>
#include <string_view>
#include <utility>
#include <vector>
Expand Down Expand Up @@ -51,34 +52,21 @@ int64_t ToPosition(int32_t key, uint32_t pos32) {
return (int64_t{key} << 32) | int64_t{pos32};
}

Status ValidatePosition(int64_t pos) {
if (pos < 0 || pos > RoaringPositionBitmap::kMaxPosition) {
return InvalidArgument("Bitmap supports positions that are >= 0 and <= {}: {}",
RoaringPositionBitmap::kMaxPosition, pos);
}
return {};
}

void WriteBitmaps(const std::vector<roaring::Roaring>& bitmaps, uint8_t* buf) {
void WriteBitmaps(const std::map<int32_t, roaring::Roaring>& bitmaps, uint8_t* buf) {
WriteLittleEndian(static_cast<int64_t>(bitmaps.size()), buf);
buf += kBitmapCountSizeBytes;
for (int32_t key = 0; std::cmp_less(key, bitmaps.size()); ++key) {
for (const auto& [key, bitmap] : bitmaps) {
WriteLittleEndian(key, buf);
buf += kBitmapKeySizeBytes;
buf += bitmaps[key].write(reinterpret_cast<char*>(buf), /*portable=*/true);
buf += bitmap.write(reinterpret_cast<char*>(buf), /*portable=*/true);
}
}

} // namespace

struct RoaringPositionBitmap::Impl {
std::vector<roaring::Roaring> bitmaps;

void AllocateBitmapsIfNeeded(int32_t required_length) {
if (std::cmp_less(bitmaps.size(), required_length)) {
bitmaps.resize(static_cast<size_t>(required_length));
}
}
// Empty buckets are never retained.
std::map<int32_t, roaring::Roaring> bitmaps;
};

RoaringPositionBitmap::RoaringPositionBitmap() : impl_(std::make_unique<Impl>()) {}
Expand Down Expand Up @@ -108,86 +96,85 @@ RoaringPositionBitmap::RoaringPositionBitmap(std::unique_ptr<Impl> impl)
: impl_(std::move(impl)) {}

void RoaringPositionBitmap::Add(int64_t pos) {
if (pos < 0 || pos > kMaxPosition) {
if (pos < 0) {
return; // Silently ignore invalid positions
}
int32_t key = Key(pos);
uint32_t pos32 = Pos32Bits(pos);
impl_->AllocateBitmapsIfNeeded(key + 1);
impl_->bitmaps[key].add(pos32);
}

void RoaringPositionBitmap::AddManyForKey(int32_t key,
std::span<const uint32_t> positions) {
impl_->AllocateBitmapsIfNeeded(key + 1);
if (key < 0 || positions.empty()) {
return;
}
impl_->bitmaps[key].addMany(positions.size(), positions.data());
}

void RoaringPositionBitmap::AddRange(int64_t pos_start, int64_t pos_end) {
pos_start = std::max(pos_start, int64_t{0});
pos_end = std::min(pos_end, kMaxPosition + 1);
if (pos_start >= pos_end) {
return;
}

int64_t pos_last = pos_end - 1;
int32_t start_key = Key(pos_start);
int32_t end_key = Key(pos_last);
impl_->AllocateBitmapsIfNeeded(end_key + 1);

for (int32_t key = start_key; key <= end_key; ++key) {
for (int64_t key = start_key; key <= end_key; ++key) {
uint64_t low_start = (key == start_key) ? Pos32Bits(pos_start) : uint64_t{0};
uint64_t low_end = (key == end_key) ? static_cast<uint64_t>(Pos32Bits(pos_last)) + 1
: (uint64_t{1} << 32);
impl_->bitmaps[key].addRange(low_start, low_end);
impl_->bitmaps[static_cast<int32_t>(key)].addRange(low_start, low_end);
}
}

bool RoaringPositionBitmap::Contains(int64_t pos) const {
if (pos < 0 || pos > kMaxPosition) {
if (pos < 0) {
return false; // Invalid positions are not contained
}
int32_t key = Key(pos);
uint32_t pos32 = Pos32Bits(pos);
return std::cmp_less(key, impl_->bitmaps.size()) && impl_->bitmaps[key].contains(pos32);
auto it = impl_->bitmaps.find(key);
return it != impl_->bitmaps.end() && it->second.contains(pos32);
}

bool RoaringPositionBitmap::IsEmpty() const { return Cardinality() == 0; }
bool RoaringPositionBitmap::IsEmpty() const { return impl_->bitmaps.empty(); }

size_t RoaringPositionBitmap::Cardinality() const {
size_t total = 0;
for (const auto& bitmap : impl_->bitmaps) {
for (const auto& [_, bitmap] : impl_->bitmaps) {
total += bitmap.cardinality();
}
return total;
}

void RoaringPositionBitmap::Or(const RoaringPositionBitmap& other) {
impl_->AllocateBitmapsIfNeeded(static_cast<int32_t>(other.impl_->bitmaps.size()));
for (size_t key = 0; key < other.impl_->bitmaps.size(); ++key) {
impl_->bitmaps[key] |= other.impl_->bitmaps[key];
for (const auto& [key, bitmap] : other.impl_->bitmaps) {
impl_->bitmaps[key] |= bitmap;
}
}

bool RoaringPositionBitmap::Optimize() {
bool changed = false;
for (auto& bitmap : impl_->bitmaps) {
for (auto& [_, bitmap] : impl_->bitmaps) {
changed |= bitmap.runOptimize();
}
return changed;
}

void RoaringPositionBitmap::ForEach(const std::function<void(int64_t)>& fn) const {
for (size_t key = 0; key < impl_->bitmaps.size(); ++key) {
for (uint32_t pos32 : impl_->bitmaps[key]) {
fn(ToPosition(static_cast<int32_t>(key), pos32));
for (const auto& [key, bitmap] : impl_->bitmaps) {
for (uint32_t pos32 : bitmap) {
fn(ToPosition(key, pos32));
}
}
}

size_t RoaringPositionBitmap::SerializedSizeInBytes() const {
size_t size = kBitmapCountSizeBytes;
for (const auto& bitmap : impl_->bitmaps) {
for (const auto& [_, bitmap] : impl_->bitmaps) {
size += kBitmapKeySizeBytes + bitmap.getSizeInBytes(/*portable=*/true);
}
return size;
Expand Down Expand Up @@ -238,18 +225,10 @@ Result<RoaringPositionBitmap> RoaringPositionBitmap::Deserialize(std::string_vie
remaining -= kBitmapKeySizeBytes;

ICEBERG_PRECHECK(key >= 0, "Invalid unsigned key: {}", key);
ICEBERG_PRECHECK(key < std::numeric_limits<int32_t>::max(), "Key is too large: {}",
key);
ICEBERG_PRECHECK(key > last_key,
"Keys must be sorted in ascending order, got key {} after {}", key,
last_key);

// Fill gaps with empty bitmaps
while (last_key < key - 1) {
impl->bitmaps.emplace_back();
++last_key;
}

// Read bitmap using portable safe deserialization.
// CRoaring's readSafe may throw on corrupted data.
roaring::Roaring bitmap;
Expand All @@ -266,7 +245,9 @@ Result<RoaringPositionBitmap> RoaringPositionBitmap::Deserialize(std::string_vie
buf += bitmap_size;
remaining -= bitmap_size;

impl->bitmaps.emplace_back(std::move(bitmap));
if (!bitmap.isEmpty()) {
impl->bitmaps.emplace(key, std::move(bitmap));
}
last_key = key;
--remaining_count;
}
Expand Down
23 changes: 12 additions & 11 deletions src/iceberg/deletes/roaring_position_bitmap.h
Original file line number Diff line number Diff line change
Expand Up @@ -20,10 +20,11 @@
#pragma once

/// \file iceberg/deletes/roaring_position_bitmap.h
/// A 64-bit position bitmap using an array of 32-bit Roaring bitmaps.
/// A 64-bit position bitmap using sparse 32-bit Roaring bitmaps.

#include <cstdint>
#include <functional>
#include <limits>
#include <memory>
#include <span>
#include <string>
Expand All @@ -37,7 +38,7 @@ namespace iceberg {

class PositionDeleteIndex;

/// \brief A bitmap that supports positive 64-bit positions, optimized
/// \brief A bitmap that supports non-negative 64-bit positions, optimized
/// for cases where most positions fit in 32 bits.
///
/// Incoming 64-bit positions are divided into a 32-bit "key" using the
Expand All @@ -51,8 +52,8 @@ class PositionDeleteIndex;
/// for `deletion-vector-v1` persistence.
class ICEBERG_EXPORT RoaringPositionBitmap {
public:
/// \brief Maximum supported position (aligned with the Java implementation).
static constexpr int64_t kMaxPosition = 0x7FFFFFFE80000000LL;
/// \brief Maximum supported position.
static constexpr int64_t kMaxPosition = std::numeric_limits<int64_t>::max();

RoaringPositionBitmap();
~RoaringPositionBitmap();
Expand All @@ -64,21 +65,21 @@ class ICEBERG_EXPORT RoaringPositionBitmap {
RoaringPositionBitmap& operator=(const RoaringPositionBitmap& other);

/// \brief Sets a position in the bitmap.
/// \param pos the position (must be >= 0 and <= kMaxPosition)
/// \note Invalid positions are silently ignored
/// \param pos the position (must be non-negative)
/// \note Negative positions are silently ignored.
void Add(int64_t pos);

/// \brief Sets a range of positions [pos_start, pos_end).
/// \param pos_start the start of the range (inclusive), clamped to 0
/// \param pos_end the end of the range (exclusive), clamped to kMaxPosition + 1
/// \note If pos_start > pos_end, the call is silently ignored.
/// If pos_start == pos_end, this method does nothing.
/// Positions outside [0, kMaxPosition] are silently ignored.
/// \param pos_end the end of the range (exclusive)
/// \note Empty and reversed ranges are silently ignored.
/// \note Because pos_end is an int64_t exclusive endpoint, this method cannot
/// include kMaxPosition. Call Add(kMaxPosition) separately.
void AddRange(int64_t pos_start, int64_t pos_end);

/// \brief Checks if a position is set in the bitmap.
/// \param pos the position to check
/// \return true if the position is set, false otherwise (including invalid positions)
/// \return true if the position is set, false otherwise (including negative positions)
bool Contains(int64_t pos) const;

/// \brief Returns true if the bitmap has no positions set.
Expand Down
Loading