Skip to content
Merged
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
88 changes: 80 additions & 8 deletions src/paimon/core/manifest/manifest_file.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,10 @@
#include "arrow/c/abi.h"
#include "arrow/c/bridge.h"
#include "paimon/common/data/columnar/columnar_row.h"
#include "paimon/common/predicate/predicate_validator.h"
#include "paimon/common/reader/late_materializing_file_batch_reader.h"
#include "paimon/common/utils/arrow/status_utils.h"
#include "paimon/common/utils/checked_cast.h"
#include "paimon/core/io/rolling_file_writer.h"
#include "paimon/core/manifest/manifest_entry.h"
#include "paimon/core/manifest/manifest_entry_serializer.h"
Expand All @@ -39,6 +42,7 @@
#include "paimon/format/file_format.h"
#include "paimon/format/reader_builder.h"
#include "paimon/format/writer_builder.h"
#include "paimon/predicate/predicate_builder.h"
#include "paimon/status.h"

namespace arrow {
Expand All @@ -49,6 +53,12 @@ class Schema;
namespace paimon {
class MemoryPool;

namespace {
constexpr int32_t kVersionFieldIndex = 0;
constexpr int32_t kBucketFieldIndex = 3;
constexpr int32_t kTotalBucketsFieldIndex = 4;
} // namespace

ManifestFile::ManifestFile(const std::shared_ptr<FileSystem>& file_system,
const std::shared_ptr<ReaderBuilder>& reader_builder,
const std::shared_ptr<WriterBuilder>& writer_builder,
Expand Down Expand Up @@ -91,25 +101,87 @@ Result<std::unique_ptr<ManifestFile>> ManifestFile::Create(
}

Status ManifestFile::ReadBucketEntries(const std::string& file_name, int32_t bucket,
const std::optional<int32_t>& expected_total_buckets,
std::optional<int64_t> file_size,
std::vector<ManifestEntry>* entries) const {
// Readers without precise bitmap selection still filter aligned Arrow columns
// before constructing ManifestEntry and DataFileMeta objects.
Comment thread
lxy-9602 marked this conversation as resolved.
return ReadArrowBatches(
file_name,
[this, bucket, entries](const std::shared_ptr<arrow::StructArray>& batch) -> Status {
const arrow::ArrayVector& fields = batch->fields();
ColumnarRow row(fields, pool_, /*row_id=*/0);
for (int64_t i = 0; i < batch->length(); i++) {
file_name, file_size,
[this, bucket, expected_total_buckets,
entries](const std::shared_ptr<arrow::StructArray>& batch) -> Status {
ColumnarRow row(batch->fields(), pool_, /*row_id=*/0);
for (int64_t i = 0; i < batch->length(); ++i) {
row.SetRowId(i);
PAIMON_RETURN_NOT_OK(ManifestEntrySerializer::ValidateVersion(row.GetInt(0)));
if (ManifestEntrySerializer::GetBucket(row) != bucket) {
PAIMON_RETURN_NOT_OK(
ManifestEntrySerializer::ValidateVersion(row.GetInt(kVersionFieldIndex)));
// Different or unknown bucket counts must reach per-entry bucket filtering.
const bool historical_layout =
expected_total_buckets &&
(row.IsNullAt(kBucketFieldIndex) || row.IsNullAt(kTotalBucketsFieldIndex) ||
row.GetInt(kTotalBucketsFieldIndex) != expected_total_buckets.value());
if (!historical_layout && ManifestEntrySerializer::GetBucket(row) != bucket) {
continue;
}
PAIMON_ASSIGN_OR_RAISE(ManifestEntry entry, serializer_->FromRow(row));
entries->push_back(std::move(entry));
}
return Status::OK();
},
file_size);
[this, bucket, expected_total_buckets](std::unique_ptr<FileBatchReader>* reader) {
return PrepareBucketRead(bucket, expected_total_buckets, reader);
});
}

Status ManifestFile::PrepareBucketRead(int32_t bucket,
const std::optional<int32_t>& expected_total_buckets,
std::unique_ptr<FileBatchReader>* reader) const {
if (!(*reader)->SupportPreciseBitmapSelection()) {
Comment thread
lxy-9602 marked this conversation as resolved.
return Status::OK();
}
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<ArrowSchema> c_schema, (*reader)->GetFileSchema());
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Schema> file_schema,
arrow::ImportSchema(c_schema.get()));
const auto& target_type = serializer_->GetDataType();
const std::string& bucket_name = target_type->field(kBucketFieldIndex)->name();
std::shared_ptr<Predicate> selector =
PredicateBuilder::Equal(kBucketFieldIndex, bucket_name, FieldType::INT, Literal(bucket));
if (expected_total_buckets) {
const std::string& total_name = target_type->field(kTotalBucketsFieldIndex)->name();
PAIMON_ASSIGN_OR_RAISE(
selector,
PredicateBuilder::Or(
{selector,
PredicateBuilder::IsNull(kTotalBucketsFieldIndex, total_name, FieldType::INT),
PredicateBuilder::NotEqual(kTotalBucketsFieldIndex, total_name, FieldType::INT,
Literal(expected_total_buckets.value()))}));
}
// Retain unsupported versions regardless of bucket so the consumer validates every version
// before bucket filtering, including when the probe would otherwise select no entries.
const std::string& version_name = target_type->field(kVersionFieldIndex)->name();
PAIMON_ASSIGN_OR_RAISE(
selector,
Comment thread
lxy-9602 marked this conversation as resolved.
PredicateBuilder::Or(
{selector, PredicateBuilder::IsNull(kVersionFieldIndex, version_name, FieldType::INT),
PredicateBuilder::NotEqual(
kVersionFieldIndex, version_name, FieldType::INT,
Literal(
checked_cast<ManifestEntrySerializer*>(serializer_.get())->GetVersion()))}));
if (!PredicateValidator::ValidatePredicateWithSchema(*file_schema, selector,
/*validate_field_idx=*/true)
.ok()) {
// An incompatible projection must fall back to ordinary schema evolution.
return Status::OK();
}
PAIMON_ASSIGN_OR_RAISE(
std::unique_ptr<LateMaterializingFileBatchReader> selective_reader,
LateMaterializingFileBatchReader::Create(std::move(*reader), arrow_pool_));
// Keep the on-disk schema; ManifestMetaReader still performs schema evolution afterwards.
ArrowSchema full_c_schema;
PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*file_schema, &full_c_schema));
PAIMON_RETURN_NOT_OK(selective_reader->SetReadSchema(&full_c_schema, selector, std::nullopt));
*reader = std::move(selective_reader);
return Status::OK();
}

Result<std::vector<ManifestFileMeta>> ManifestFile::Write(
Expand Down
9 changes: 8 additions & 1 deletion src/paimon/core/manifest/manifest_file.h
Original file line number Diff line number Diff line change
Expand Up @@ -63,16 +63,23 @@ class ManifestFile : public ObjectsFile<ManifestEntry> {
/// @note This method is atomic.
Result<std::vector<ManifestFileMeta>> Write(const std::vector<ManifestEntry>& entries);

/// Read a manifest file and deserialize only entries for the specified bucket.
/// Read entries for a bucket. An inferred total bucket count also retains entries with
/// different or unknown bucket counts for per-entry bucket filtering.
/// Bucket pruning assumes stable bucket-key hashing across schema versions.
/// Serialization versions are validated for all entries before bucket filtering.
///
/// @param file_size Length of the manifest when the caller already has it from the manifest
/// list, which saves the read a metadata request on a remote store. Pass
/// std::nullopt when the length is not known.
Status ReadBucketEntries(const std::string& file_name, int32_t bucket,
const std::optional<int32_t>& expected_total_buckets,
std::optional<int64_t> file_size,
std::vector<ManifestEntry>* entries) const;

private:
Status PrepareBucketRead(int32_t bucket, const std::optional<int32_t>& expected_total_buckets,
std::unique_ptr<FileBatchReader>* reader) const;

ManifestFile(const std::shared_ptr<FileSystem>& file_system,
const std::shared_ptr<ReaderBuilder>& reader_builder,
const std::shared_ptr<WriterBuilder>& writer_builder,
Expand Down
Loading
Loading