lxy-9602 commented on code in PR #323:
URL: https://github.com/apache/paimon-cpp/pull/323#discussion_r4001667497
##########
src/paimon/core/operation/file_store_scan.cpp:
##########
@@ -464,6 +482,24 @@ Status FileStoreScan::ReadAndMergeBucketFileEntries(
return MergeLiveEntries(unmerged_entries, merged_entries);
}
+Result<bool> FileStoreScan::CheckHistoricalBucketCompatibility(
+ const std::vector<ManifestFileMeta>& manifest_metas) const {
+ int64_t max_schema_id = -1;
+ for (const auto& meta : manifest_metas) {
+ max_schema_id = std::max(max_schema_id, meta.SchemaId());
+ }
+ for (int64_t schema_id = 0; schema_id <= max_schema_id; ++schema_id) {
+ PAIMON_ASSIGN_OR_RAISE(bool compatible,
HasCompatibleBucketKeys(schema_id));
+ if (!compatible) {
+ return false;
+ }
+ if (schema_id == max_schema_id) {
+ break;
+ }
+ }
+ return true;
+}
Review Comment:
The current implementation scans all schemas, which could have a significant
performance impact. In real-world environments, the schema may change many
times, and some of those schema versions may not even produce a snapshot. Also,
Java’s `BucketFilter` only considers (partition, bucket, totalBuckets) and does
not check the data schema ID.
##########
src/paimon/core/manifest/manifest_file.cpp:
##########
@@ -90,25 +92,79 @@ Result<std::unique_ptr<ManifestFile>> ManifestFile::Create(
}
Status ManifestFile::ReadBucketEntries(const std::string& file_name, int32_t
bucket,
- std::vector<ManifestEntry>* entries)
const {
+ std::vector<ManifestEntry>* entries,
+ const std::optional<int32_t>&
expected_total_buckets) const {
+ // Readers without an in-memory selective probe still filter aligned Arrow
columns
+ // before constructing ManifestEntry and DataFileMeta objects.
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++) {
+ [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) {
+ // Different or unknown bucket counts must reach the
compatibility checks.
+ const bool historical_layout =
+ expected_total_buckets && (row.IsNullAt(3) ||
row.IsNullAt(4) ||
+ row.GetInt(4) !=
expected_total_buckets.value());
+ if (!historical_layout &&
ManifestEntrySerializer::GetBucket(row) != bucket) {
continue;
}
+ // Only validate entries retained by the bucket selector.
FromRow checks the
+ // serialization version before decoding the remaining
metadata.
PAIMON_ASSIGN_OR_RAISE(ManifestEntry entry,
serializer_->FromRow(row));
entries->push_back(std::move(entry));
}
return Status::OK();
+ },
+ [this, bucket,
expected_total_buckets](std::unique_ptr<FileBatchReader>* reader) {
+ return PrepareBucketRead(reader, bucket, expected_total_buckets);
});
}
+Status ManifestFile::PrepareBucketRead(std::unique_ptr<FileBatchReader>*
reader, int32_t bucket,
+ const std::optional<int32_t>&
expected_total_buckets) const {
+ if (!(*reader)->SupportPreciseBitmapSelection()) {
+ 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(3)->name();
+ auto bucket_field = file_schema->GetFieldByName(bucket_name);
+ if (!bucket_field || bucket_field->type()->id() != arrow::Type::INT32) {
+ return Status::OK();
+ }
+ if (expected_total_buckets) {
+ auto total_field = file_schema->GetFieldByName("_TOTAL_BUCKETS");
+ if (!total_field || total_field->type()->id() != arrow::Type::INT32) {
+ return Status::OK();
+ }
+ }
Review Comment:
Instead of manually validating each field type, could we use
`PredicateValidator::ValidatePredicateWithSchema()` here for consistency?
##########
src/paimon/core/operation/file_store_scan.cpp:
##########
@@ -422,15 +422,33 @@ Status FileStoreScan::ReadAndMergeBucketFileEntries(
const std::vector<ManifestFileMeta>& manifest_metas, int32_t bucket,
std::vector<ManifestEntry>* merged_entries) const {
const bool inferred_bucket = !bucket_filter_ && bucket_selector_ !=
nullptr;
- // Explicit-bucket lazy decoding cannot retain entries with a different
layout.
- if (!inferred_bucket &&
core_options_.ScanManifestEntryLazyDecodeEnabled()) {
+ bool use_lazy_decode = core_options_.ScanManifestEntryLazyDecodeEnabled();
+ if (use_lazy_decode && inferred_bucket) {
+ // A manifest's schema ID is only an upper bound. Prove the whole
historical range
+ // compatible before pruning, and fall back if an unused historical
schema cannot be read.
+ Result<bool> compatible =
CheckHistoricalBucketCompatibility(manifest_metas);
+ use_lazy_decode = compatible.ok() && compatible.value();
+ }
+ if (use_lazy_decode) {
std::vector<std::future<Result<std::vector<ManifestEntry>>>> futures;
futures.reserve(manifest_metas.size());
for (const auto& meta : manifest_metas) {
- auto read_meta_task = [this, meta, bucket]() ->
Result<std::vector<ManifestEntry>> {
+ auto read_meta_task = [this, meta, bucket,
+ inferred_bucket]() ->
Result<std::vector<ManifestEntry>> {
std::vector<ManifestEntry> bucket_entries;
- PAIMON_RETURN_NOT_OK(
- manifest_file_->ReadBucketEntries(meta.FileName(), bucket,
&bucket_entries));
+ if (inferred_bucket) {
+ PAIMON_RETURN_NOT_OK(manifest_file_->ReadBucketEntries(
+ meta.FileName(), bucket, &bucket_entries,
core_options_.GetBucket()));
+ } else if (meta.MinBucket() && meta.MaxBucket() &&
+ meta.MinBucket().value() == bucket &&
+ meta.MaxBucket().value() == bucket) {
+ // Every entry belongs to this bucket; a projection pass
cannot prune rows.
+ PAIMON_RETURN_NOT_OK(
+ manifest_file_->Read(meta.FileName(),
/*filter=*/nullptr, &bucket_entries));
Review Comment:
I don’t fully understand the logic here. If `min == max == bucket`, it seems
we wouldn’t need late filtering at all, and applying it would only increase I/O.
##########
src/paimon/core/utils/objects_file.h:
##########
@@ -201,6 +206,9 @@ Status ObjectsFile<T>::ReadArrowBatches(
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<FileBatchReader> batch_reader,
reader_builder_->Build(file_input_stream));
+ if (prepare_reader && cached_bytes) {
+ PAIMON_RETURN_NOT_OK(prepare_reader(&batch_reader));
+ }
Review Comment:
It seems that Java performs late materialization for partition and bucket
regardless of whether caching is enabled. I’m not sure what the benefit is of
tying late materialization to cache in the current implementation.
Alternatively, please provide test results showing whether enabling late
materialization without cache has any negative impact.
##########
src/paimon/core/manifest/manifest_file.cpp:
##########
@@ -90,25 +92,79 @@ Result<std::unique_ptr<ManifestFile>> ManifestFile::Create(
}
Status ManifestFile::ReadBucketEntries(const std::string& file_name, int32_t
bucket,
- std::vector<ManifestEntry>* entries)
const {
+ std::vector<ManifestEntry>* entries,
+ const std::optional<int32_t>&
expected_total_buckets) const {
+ // Readers without an in-memory selective probe still filter aligned Arrow
columns
+ // before constructing ManifestEntry and DataFileMeta objects.
Review Comment:
Similarly, `entries` is an output parameter, and output parameters should
generally be placed at the end of the parameter list. Please don’t insert the
new parameter after it.
##########
src/paimon/core/manifest/manifest_file.cpp:
##########
@@ -90,25 +92,79 @@ Result<std::unique_ptr<ManifestFile>> ManifestFile::Create(
}
Status ManifestFile::ReadBucketEntries(const std::string& file_name, int32_t
bucket,
- std::vector<ManifestEntry>* entries)
const {
+ std::vector<ManifestEntry>* entries,
+ const std::optional<int32_t>&
expected_total_buckets) const {
+ // Readers without an in-memory selective probe still filter aligned Arrow
columns
+ // before constructing ManifestEntry and DataFileMeta objects.
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++) {
+ [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) {
+ // Different or unknown bucket counts must reach the
compatibility checks.
+ const bool historical_layout =
+ expected_total_buckets && (row.IsNullAt(3) ||
row.IsNullAt(4) ||
+ row.GetInt(4) !=
expected_total_buckets.value());
+ if (!historical_layout &&
ManifestEntrySerializer::GetBucket(row) != bucket) {
continue;
}
+ // Only validate entries retained by the bucket selector.
FromRow checks the
+ // serialization version before decoding the remaining
metadata.
PAIMON_ASSIGN_OR_RAISE(ManifestEntry entry,
serializer_->FromRow(row));
entries->push_back(std::move(entry));
}
return Status::OK();
+ },
+ [this, bucket,
expected_total_buckets](std::unique_ptr<FileBatchReader>* reader) {
+ return PrepareBucketRead(reader, bucket, expected_total_buckets);
});
}
+Status ManifestFile::PrepareBucketRead(std::unique_ptr<FileBatchReader>*
reader, int32_t bucket,
+ const std::optional<int32_t>&
expected_total_buckets) const {
+ if (!(*reader)->SupportPreciseBitmapSelection()) {
Review Comment:
It looks like `reader` is an output parameter. Please place output
parameters at the end of the function parameter list.
##########
src/paimon/core/manifest/manifest_file.cpp:
##########
@@ -90,25 +92,79 @@ Result<std::unique_ptr<ManifestFile>> ManifestFile::Create(
}
Status ManifestFile::ReadBucketEntries(const std::string& file_name, int32_t
bucket,
- std::vector<ManifestEntry>* entries)
const {
+ std::vector<ManifestEntry>* entries,
+ const std::optional<int32_t>&
expected_total_buckets) const {
+ // Readers without an in-memory selective probe still filter aligned Arrow
columns
+ // before constructing ManifestEntry and DataFileMeta objects.
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++) {
+ [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) {
+ // Different or unknown bucket counts must reach the
compatibility checks.
Review Comment:
I’m a bit curious why `ValidateVersion` was removed. It seems like the
result should still be version-validated even after late materialization.
##########
src/paimon/core/manifest/manifest_file.cpp:
##########
@@ -90,25 +92,79 @@ Result<std::unique_ptr<ManifestFile>> ManifestFile::Create(
}
Status ManifestFile::ReadBucketEntries(const std::string& file_name, int32_t
bucket,
- std::vector<ManifestEntry>* entries)
const {
+ std::vector<ManifestEntry>* entries,
+ const std::optional<int32_t>&
expected_total_buckets) const {
+ // Readers without an in-memory selective probe still filter aligned Arrow
columns
+ // before constructing ManifestEntry and DataFileMeta objects.
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++) {
+ [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) {
+ // Different or unknown bucket counts must reach the
compatibility checks.
+ const bool historical_layout =
+ expected_total_buckets && (row.IsNullAt(3) ||
row.IsNullAt(4) ||
+ row.GetInt(4) !=
expected_total_buckets.value());
+ if (!historical_layout &&
ManifestEntrySerializer::GetBucket(row) != bucket) {
continue;
}
+ // Only validate entries retained by the bucket selector.
FromRow checks the
+ // serialization version before decoding the remaining
metadata.
PAIMON_ASSIGN_OR_RAISE(ManifestEntry entry,
serializer_->FromRow(row));
entries->push_back(std::move(entry));
}
return Status::OK();
+ },
+ [this, bucket,
expected_total_buckets](std::unique_ptr<FileBatchReader>* reader) {
+ return PrepareBucketRead(reader, bucket, expected_total_buckets);
});
}
+Status ManifestFile::PrepareBucketRead(std::unique_ptr<FileBatchReader>*
reader, int32_t bucket,
+ const std::optional<int32_t>&
expected_total_buckets) const {
+ if (!(*reader)->SupportPreciseBitmapSelection()) {
+ 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(3)->name();
Review Comment:
There are too many hard-coded indices like `3` and `4` here. Please replace
them with descriptively named constants instead.
##########
src/paimon/format/avro/avro_file_batch_reader.cpp:
##########
@@ -180,6 +235,24 @@ Status AvroFileBatchReader::SetReadSchema(::ArrowSchema*
read_schema,
reader_ = std::move(reader);
array_builder_ = std::move(array_builder);
decode_context_.ClearBuilderMetadata();
+ selection_iterator_.reset();
+ selection_end_.reset();
+ selection_bitmap_ = selection_bitmap;
+ if (!block_index_complete_) {
+ block_index_.clear();
+ block_index_disabled_ = false;
+ }
+ if (selection_bitmap_ && !selection_bitmap_->IsEmpty() &&
block_index_complete_ &&
+ !block_index_.empty()) {
+ selection_iterator_ = selection_bitmap_->Begin();
+ selection_end_ = selection_bitmap_->End();
+ const uint64_t selected_row =
static_cast<uint32_t>(**selection_iterator_);
+ auto block =
+ std::upper_bound(block_index_.begin(), block_index_.end(),
selected_row,
+ [](uint64_t row, const auto& entry) { return row
< entry.first; });
+ selected_block_ = std::distance(block_index_.begin(), block) - 1;
+ }
+ previous_row_ids_.clear();
Review Comment:
The current support for the selected bitmap in `avro_file_batch_reader.cpp`
does not seem very maintainable. Please refactor this part. I’m only listing a
few of the issues here:
- The implementation is tightly coupled with the block-optimization-based
late materialization logic, and appears to assume that the reader will first do
a full read without a bitmap and then re-read with a bitmap.
- `bitmap.Contains()` is relatively expensive. If the access pattern is
sequential, it would be better to use an iterator instead.
- This feature introduces too many combinations of member state, and the
validity of those states seems to rely entirely on implicit conventions. In
addition, ReadRowsIntoBuilder() appears to be taking on too many
responsibilities.
--
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]