wangyong9999 commented on code in PR #323:
URL: https://github.com/apache/paimon-cpp/pull/323#discussion_r3978562353


##########
src/paimon/format/avro/avro_file_batch_reader.cpp:
##########
@@ -18,23 +18,63 @@
 
 #include "paimon/format/avro/avro_file_batch_reader.h"
 
+#include <algorithm>
 #include <limits>
 #include <memory>
 #include <utility>
 
+#include "arrow/array/builder_nested.h"
 #include "arrow/c/bridge.h"
 #include "fmt/format.h"
 #include "paimon/common/metrics/metrics_impl.h"
 #include "paimon/common/utils/arrow/arrow_utils.h"
 #include "paimon/common/utils/arrow/mem_utils.h"
 #include "paimon/common/utils/arrow/status_utils.h"
+#include "paimon/common/utils/checked_cast.h"
 #include "paimon/common/utils/scope_guard.h"
 #include "paimon/core/utils/nested_projection_utils.h"
 #include "paimon/format/avro/avro_input_stream_impl.h"
 #include "paimon/format/avro/avro_schema_converter.h"
 #include "paimon/reader/batch_reader.h"
 
 namespace paimon::avro {
+namespace {
+
+// Full reads may differ in field nullability without requiring a projection. 
Compare
+// names and positions recursively so reordered fields and type changes still 
fail.
+bool SameReadLayout(const std::shared_ptr<arrow::DataType>& file_type,
+                    const std::shared_ptr<arrow::DataType>& read_type) {
+    if (file_type->id() != read_type->id()) {
+        return false;
+    }
+    switch (file_type->id()) {
+        case arrow::Type::MAP:
+            if (checked_pointer_cast<arrow::MapType>(file_type)->keys_sorted() 
!=
+                
checked_pointer_cast<arrow::MapType>(read_type)->keys_sorted()) {
+                return false;
+            }
+            break;
+        case arrow::Type::STRUCT:
+        case arrow::Type::LIST:
+            break;
+        default:
+            return file_type->Equals(read_type);

Review Comment:
   This rejects nested types that Avro can't represent exactly. 
TINYINT/SMALLINT are written as Avro `int`, so the file type here is `int32` 
while the data-schema read type is `int8`/`int16`; the decoder already handles 
that pair through its builder-type dispatch. Before this PR a `ROW<a TINYINT>` 
column in an Avro data file read fine, now `SetReadSchema` fails with "Avro 
full struct read requires matching field names, order and types" (reproduced 
locally with `struct<r: struct<a: int8, b: int16>>` against a file written from 
the same type). Top-level columns aren't type-checked at all, so 
`ARRAY<TINYINT>` keeps working while `ROW<a TINYINT>` breaks.
   
   Suggest comparing only names, order and nesting kind here (that's all 
positional decoding needs) and leaving leaf type compatibility to the decoder, 
as the top-level path already does.
   



##########
src/paimon/core/operation/file_store_scan.cpp:
##########
@@ -422,15 +422,27 @@ 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()) {
+    if (core_options_.ScanManifestEntryLazyDecodeEnabled()) {
         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_->ReadInferredBucketEntries(
+                        meta.FileName(), bucket, core_options_.GetBucket(), 
table_schema_->Id(),
+                        &bucket_entries));
+                } 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));
+                } else {
+                    
PAIMON_RETURN_NOT_OK(manifest_file_->ReadBucketEntries(meta.FileName(), bucket,

Review Comment:
   Lazy decode is on by default, so every cached explicit-bucket read now takes 
the probe pass regardless of how many buckets the manifest spans. The profile 
you posted is the 16K-bucket case; for the usual few-bucket tables this is a 
second full parse plus a second zstd decompression of nearly every block, to 
save materializing only a fraction of the rows, which can easily be a net loss. 
`meta.MinBucket()/MaxBucket()` is already available here, so gating the probe 
on the bucket span would keep the win for the wide case without taxing the 
common one.
   



##########
src/paimon/core/manifest/manifest_file.cpp:
##########
@@ -106,9 +108,175 @@ Status ManifestFile::ReadBucketEntries(const std::string& 
file_name, int32_t buc
                 entries->push_back(std::move(entry));
             }
             return Status::OK();
+        },
+        [this, bucket](FileBatchReader* reader) { return 
PrepareBucketRead(reader, bucket); });
+}
+
+Status ManifestFile::ReadInferredBucketEntries(const std::string& file_name, 
int32_t bucket,
+                                               int32_t total_buckets, int64_t 
schema_id,
+                                               std::vector<ManifestEntry>* 
entries) 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, total_buckets, schema_id,
+         entries](const std::shared_ptr<arrow::StructArray>& batch) -> Status {
+            auto file_array = batch->GetFieldByName("_FILE");
+            if (!file_array || file_array->type_id() != arrow::Type::STRUCT) {
+                return Status::Invalid("Manifest file metadata must be a 
struct array");
+            }
+            auto files = checked_pointer_cast<arrow::StructArray>(file_array);
+            auto schema_array = files->GetFieldByName("_SCHEMA_ID");
+            std::shared_ptr<arrow::Int64Array> schema_ids;
+            if (schema_array && schema_array->type_id() == arrow::Type::INT64) 
{
+                schema_ids = 
checked_pointer_cast<arrow::Int64Array>(schema_array);
+            }
+            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)));
+                // Only discard other buckets when both layout identifiers 
match. Unknown
+                // or historical layouts must reach the existing compatibility 
checks.
+                if (!row.IsNullAt(3) && 
ManifestEntrySerializer::GetBucket(row) != bucket &&
+                    !row.IsNullAt(4) && row.GetInt(4) == total_buckets && 
!files->IsNull(i) &&
+                    schema_ids && !schema_ids->IsNull(i) && 
schema_ids->Value(i) == schema_id) {
+                    continue;
+                }
+                PAIMON_ASSIGN_OR_RAISE(ManifestEntry entry, 
serializer_->FromRow(row));
+                entries->push_back(std::move(entry));
+            }
+            return Status::OK();
+        },
+        [this, bucket, total_buckets, schema_id](FileBatchReader* reader) {
+            return PrepareBucketRead(reader, bucket, 
std::make_pair(total_buckets, schema_id));
         });
 }
 
+Status ManifestFile::PrepareBucketRead(
+    FileBatchReader* reader, int32_t bucket,
+    const std::optional<std::pair<int32_t, int64_t>>& inferred_layout) 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& version_name = target_type->field(0)->name();
+    const std::string& bucket_name = target_type->field(3)->name();
+    auto version_field = file_schema->GetFieldByName(version_name);
+    auto bucket_field = file_schema->GetFieldByName(bucket_name);
+    if (!version_field || !bucket_field || version_field->type()->id() != 
arrow::Type::INT32 ||
+        bucket_field->type()->id() != arrow::Type::INT32) {
+        return Status::OK();
+    }
+    std::shared_ptr<arrow::Field> projected_file;
+    if (inferred_layout) {
+        auto total_field = file_schema->GetFieldByName("_TOTAL_BUCKETS");
+        auto file_field = file_schema->GetFieldByName("_FILE");
+        if (!total_field || total_field->type()->id() != arrow::Type::INT32 || 
!file_field ||
+            file_field->type()->id() != arrow::Type::STRUCT) {
+            return Status::OK();
+        }
+        auto file_type = 
checked_pointer_cast<arrow::StructType>(file_field->type());
+        auto schema_field = file_type->GetFieldByName("_SCHEMA_ID");
+        if (!schema_field || schema_field->type()->id() != arrow::Type::INT64) 
{
+            return Status::OK();
+        }
+        projected_file = file_field->WithType(arrow::struct_({schema_field}));
+    }
+    arrow::FieldVector probe_fields;
+    for (const auto& field : file_schema->fields()) {
+        if (field->name() == version_name || field->name() == bucket_name) {
+            probe_fields.push_back(field);
+        } else if (inferred_layout && field->name() == "_TOTAL_BUCKETS") {
+            probe_fields.push_back(field);
+        } else if (inferred_layout && field->name() == "_FILE") {
+            probe_fields.push_back(projected_file);
+        }
+    }
+    auto probe_schema = arrow::schema(probe_fields);
+    ArrowSchema probe_c_schema;
+    PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*probe_schema, 
&probe_c_schema));
+    auto probe_status = reader->SetReadSchema(&probe_c_schema, 
/*predicate=*/nullptr, std::nullopt);
+    if (!probe_status.ok()) {
+        if (!inferred_layout || (!probe_status.IsInvalid() && 
!probe_status.IsNotImplemented())) {
+            return probe_status;
+        }
+        // Nested projection is optional. Restore a full single pass when 
unsupported.

Review Comment:
   The only manifest reader that reports precise bitmap selection is Avro, and 
it now accepts exactly this `_FILE._SCHEMA_ID` projection, so this fallback 
never runs. Just return `probe_status`.
   



-- 
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]

Reply via email to