lxy-9602 commented on code in PR #323:
URL: https://github.com/apache/paimon-cpp/pull/323#discussion_r4011113098
##########
src/paimon/core/manifest/manifest_file.cpp:
##########
@@ -90,25 +100,89 @@ 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::vector<ManifestEntry>* entries)
const {
+ // Readers without precise bitmap selection 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) {
+ PAIMON_RETURN_NOT_OK(
+
ManifestEntrySerializer::ValidateVersion(row.GetInt(kVersionFieldIndex)));
+ // Different or unknown bucket counts must reach the
compatibility checks.
+ 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();
+ },
+ [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()) {
+ 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(
+ file_schema->GetFieldIndex(bucket_name), bucket_name, FieldType::INT,
Literal(bucket));
+ if (expected_total_buckets) {
+ const std::string& total_name =
target_type->field(kTotalBucketsFieldIndex)->name();
+ const int32_t total_index = file_schema->GetFieldIndex(total_name);
+ PAIMON_ASSIGN_OR_RAISE(
+ selector,
+ PredicateBuilder::Or(
+ {selector, PredicateBuilder::IsNull(total_index, total_name,
FieldType::INT),
+ PredicateBuilder::NotEqual(total_index, 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();
+ const int32_t version_index = file_schema->GetFieldIndex(version_name);
+ PAIMON_ASSIGN_OR_RAISE(
+ selector,
Review Comment:
I thought the version field should already be at `kVersionFieldIndex`.
Otherwise, reading the data back and directly accessing `kVersionFieldIndex`
would also be problematic. Why does the current flow still need to remap the
index here via `file_schema->GetFieldIndex`?
##########
src/paimon/format/avro/avro_file_batch_reader.cpp:
##########
@@ -145,16 +147,116 @@ Result<BatchReader::ReadBatch>
AvroFileBatchReader::NextBatch() {
}
}
+std::optional<uint64_t> AvroFileBatchReader::SelectionCursor::NextRow() const {
+ if (next_ == end_) {
+ return std::nullopt;
+ }
+ return static_cast<uint32_t>(*next_);
+}
+
+void AvroFileBatchReader::BlockIndex::Observe(uint64_t row, int64_t
file_offset) {
+ if (state_ != State::kBuilding ||
+ (!blocks_.empty() && blocks_.back().file_offset == file_offset)) {
+ return;
+ }
+ if (blocks_.size() == kMaxBlocks) {
+ blocks_.clear();
+ state_ = State::kDisabled;
+ return;
+ }
+ blocks_.push_back({row, file_offset});
+}
+
+void AvroFileBatchReader::BlockIndex::Finish(uint64_t row_count) {
+ if (state_ == State::kBuilding) {
+ row_count_ = row_count;
+ state_ = State::kReady;
+ }
+}
+
+void AvroFileBatchReader::BlockIndex::Reset() {
+ if (state_ != State::kReady) {
+ blocks_.clear();
+ state_ = State::kBuilding;
+ }
+}
Review Comment:
For immutable files with more than 64K blocks, a full read later will
repeatedly rebuild the 64K block indexes and then disable them again. It seems
we could reset only the unfinished `kBuilding` state, while keeping `kReady`
and `kDisabled` unchanged.
##########
src/paimon/core/manifest/manifest_file.cpp:
##########
@@ -90,25 +100,89 @@ 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::vector<ManifestEntry>* entries)
const {
+ // Readers without precise bitmap selection 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) {
+ PAIMON_RETURN_NOT_OK(
+
ManifestEntrySerializer::ValidateVersion(row.GetInt(kVersionFieldIndex)));
+ // Different or unknown bucket counts must reach the
compatibility checks.
+ 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();
+ },
+ [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()) {
+ 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(
+ file_schema->GetFieldIndex(bucket_name), bucket_name, FieldType::INT,
Literal(bucket));
+ if (expected_total_buckets) {
+ const std::string& total_name =
target_type->field(kTotalBucketsFieldIndex)->name();
+ const int32_t total_index = file_schema->GetFieldIndex(total_name);
+ PAIMON_ASSIGN_OR_RAISE(
+ selector,
+ PredicateBuilder::Or(
+ {selector, PredicateBuilder::IsNull(total_index, total_name,
FieldType::INT),
+ PredicateBuilder::NotEqual(total_index, 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();
+ const int32_t version_index = file_schema->GetFieldIndex(version_name);
+ PAIMON_ASSIGN_OR_RAISE(
+ selector,
Review Comment:
Manifest format changes now only append new fields at the end, so the it
seems `FieldIndex` should remain unchanged.
##########
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:
Your target scenario seems to be an online service, where queries may share
some resources such as `TableScanResources`. However, in many other scenarios,
such as Spark, each query runs as an independent task and there is no resource
sharing, so this assumption does not always hold. The `HasCompatibleBucketKeys`
logic was originally added by @wangyong9999 — could you please confirm whether
Java has similar logic? If not, this change would also introduce additional
small I/O operations. More importantly, in the current PR, the code iterates
and reads regardless of whether the schema actually has corresponding data
files, which seems too costly.
Since the bucket key is declared immutable, it should not be changed through
ALTER TABLE in principle. If the current Java implementation still allows such
corner cases to slip through schema validation, I think we should strengthen
schema validation to prevent them there, rather than adding a lot of fallback
handling in the scan path.
--
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]