This is an automated email from the ASF dual-hosted git repository.
SteNicholas pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-cpp.git
The following commit(s) were added to refs/heads/main by this push:
new ff38c0d2 refactor(shredding): unify shredding read plans for map
shredding and variant (#235)
ff38c0d2 is described below
commit ff38c0d2188cef349c96e57b97bccaa264052a03
Author: lxy <[email protected]>
AuthorDate: Tue Aug 25 17:16:20 2026 +0800
refactor(shredding): unify shredding read plans for map shredding and
variant (#235)
---
src/paimon/CMakeLists.txt | 4 +-
.../shredding/map_shared_shredding_file_reader.h | 112 ----------
... => map_shared_shredding_read_plan_factory.cpp} | 227 +++++----------------
.../map_shared_shredding_read_plan_factory.h | 49 +++++
...ap_shared_shredding_read_plan_factory_test.cpp} | 111 +++++-----
.../data/shredding/shredding_file_reader.cpp | 2 +-
src/paimon/core/operation/abstract_split_read.cpp | 76 +++----
src/paimon/core/operation/abstract_split_read.h | 14 +-
.../core/operation/data_evolution_split_read.h | 1 -
src/paimon/core/operation/merge_file_split_read.h | 1 -
src/paimon/core/operation/raw_file_split_read.h | 2 +-
11 files changed, 183 insertions(+), 416 deletions(-)
diff --git a/src/paimon/CMakeLists.txt b/src/paimon/CMakeLists.txt
index 8b3a2953..61fd7c96 100644
--- a/src/paimon/CMakeLists.txt
+++ b/src/paimon/CMakeLists.txt
@@ -173,7 +173,7 @@ set(PAIMON_COMMON_SRCS
common/data/shredding/map_shared_shredding_batch_converter.cpp
common/data/shredding/map_shared_shredding_column_allocator.cpp
common/data/shredding/lru_map_shared_shredding_column_allocator.cpp
- common/data/shredding/map_shared_shredding_file_reader.cpp
+ common/data/shredding/map_shared_shredding_read_plan_factory.cpp
common/data/shredding/shredding_file_reader.cpp
common/utils/delta_varint_compressor.cpp
common/utils/fields_comparator.cpp
@@ -682,7 +682,7 @@ if(PAIMON_BUILD_TESTS)
common/data/shredding/sequential_map_shared_shredding_column_allocator_test.cpp
common/data/shredding/map_shared_shredding_field_dict_test.cpp
common/data/shredding/map_shared_shredding_context_test.cpp
-
common/data/shredding/map_shared_shredding_file_reader_test.cpp
+
common/data/shredding/map_shared_shredding_read_plan_factory_test.cpp
STATIC_LINK_LIBS
paimon_shared
test_utils_static
diff --git
a/src/paimon/common/data/shredding/map_shared_shredding_file_reader.h
b/src/paimon/common/data/shredding/map_shared_shredding_file_reader.h
deleted file mode 100644
index 744d09fd..00000000
--- a/src/paimon/common/data/shredding/map_shared_shredding_file_reader.h
+++ /dev/null
@@ -1,112 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing,
- * software distributed under the License is distributed on an
- * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
- * KIND, either express or implied. See the License for the
- * specific language governing permissions and limitations
- * under the License.
- */
-
-#pragma once
-
-#include <map>
-#include <memory>
-#include <optional>
-#include <string>
-#include <utility>
-#include <vector>
-
-#include "arrow/api.h"
-#include "paimon/common/data/shredding/map_shared_shredding_utils.h"
-#include "paimon/memory/memory_pool.h"
-#include "paimon/reader/file_batch_reader.h"
-
-namespace paimon {
-
-class MapFieldReadPlan {
- public:
- virtual ~MapFieldReadPlan() = default;
-
- MapFieldReadPlan(const std::shared_ptr<arrow::Field>& logical_field,
- const std::shared_ptr<arrow::Field>& physical_read_field)
- : logical_field_(logical_field),
physical_read_field_(physical_read_field) {}
-
- const std::shared_ptr<arrow::Field>& LogicalField() const {
- return logical_field_;
- }
-
- const std::shared_ptr<arrow::Field>& PhysicalReadField() const {
- return physical_read_field_;
- }
-
- virtual Result<std::shared_ptr<arrow::Array>> Materialize(
- const std::shared_ptr<arrow::Array>& physical_array,
- arrow::MemoryPool* arrow_pool) const = 0;
-
- private:
- std::shared_ptr<arrow::Field> logical_field_;
- std::shared_ptr<arrow::Field> physical_read_field_;
-};
-
-class MapFieldReadPlanFactory {
- public:
- static Result<std::unique_ptr<MapFieldReadPlan>> CreateMapReadPlan(
- const std::shared_ptr<arrow::Field>& logical_map_field,
- const MapSharedShreddingFieldMeta& meta);
-
- static Result<std::unique_ptr<MapFieldReadPlan>>
CreateSharedSelectedKeysReadPlan(
- const std::shared_ptr<arrow::Field>& selected_keys_field,
- const MapSharedShreddingFieldMeta& meta);
-
- static Result<std::unique_ptr<MapFieldReadPlan>>
CreateDefaultSelectedKeysReadPlan(
- const std::shared_ptr<arrow::Field>& file_map_field,
- const std::shared_ptr<arrow::Field>& selected_keys_field);
-};
-
-class MapSharedShreddingFileReader : public FileBatchReader {
- public:
- MapSharedShreddingFileReader(
- std::unique_ptr<FileBatchReader>&& reader,
- std::map<std::string, std::unique_ptr<MapFieldReadPlan>>&&
field_read_plans,
- const std::shared_ptr<MemoryPool>& pool);
-
- Result<std::unique_ptr<::ArrowSchema>> GetFileSchema() const override;
-
- Status SetReadSchema(::ArrowSchema* read_schema, const
std::shared_ptr<Predicate>& predicate,
- const std::optional<RoaringBitmap32>&
selection_bitmap) override;
-
- Result<ReadBatch> NextBatch() override;
-
- Result<ReadBatchWithBitmap> NextBatchWithBitmap() override;
-
- std::shared_ptr<Metrics> GetReaderMetrics() const override;
-
- void Close() override;
-
- Result<uint64_t> GetPreviousBatchFileRowId(uint64_t batch_row_id) const
override;
-
- Result<uint64_t> GetNumberOfRows() const override;
-
- bool SupportPreciseBitmapSelection() const override;
-
- private:
- static Result<std::shared_ptr<arrow::Field>> ToLogicalMapField(
- const std::shared_ptr<arrow::Field>& physical_field);
-
- private:
- std::shared_ptr<arrow::MemoryPool> arrow_pool_;
- std::unique_ptr<FileBatchReader> reader_;
- std::map<std::string, std::unique_ptr<MapFieldReadPlan>> field_read_plans_;
-};
-
-} // namespace paimon
diff --git
a/src/paimon/common/data/shredding/map_shared_shredding_file_reader.cpp
b/src/paimon/common/data/shredding/map_shared_shredding_read_plan_factory.cpp
similarity index 77%
rename from
src/paimon/common/data/shredding/map_shared_shredding_file_reader.cpp
rename to
src/paimon/common/data/shredding/map_shared_shredding_read_plan_factory.cpp
index b1cdf963..4cb28f69 100644
--- a/src/paimon/common/data/shredding/map_shared_shredding_file_reader.cpp
+++
b/src/paimon/common/data/shredding/map_shared_shredding_read_plan_factory.cpp
@@ -17,19 +17,14 @@
* under the License.
*/
-#include "paimon/common/data/shredding/map_shared_shredding_file_reader.h"
+#include
"paimon/common/data/shredding/map_shared_shredding_read_plan_factory.h"
-#include <optional>
#include <set>
#include <string_view>
#include <utility>
#include <vector>
-#include "arrow/c/bridge.h"
-#include "arrow/util/key_value_metadata.h"
#include "fmt/format.h"
-#include "paimon/common/reader/reader_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/core/utils/nested_projection_utils.h"
@@ -69,16 +64,35 @@ void CollectPhysicalColumns(
}
}
-class FullMapReadPlan : public MapFieldReadPlan {
+class MapShreddingColumnReadPlan : public ShreddingColumnReadPlan {
+ public:
+ MapShreddingColumnReadPlan(std::shared_ptr<arrow::Field> logical_field,
+ std::shared_ptr<arrow::Field> physical_field)
+ : logical_field_(std::move(logical_field)),
physical_field_(std::move(physical_field)) {}
+
+ const std::shared_ptr<arrow::Field>& LogicalField() const override {
+ return logical_field_;
+ }
+
+ const std::shared_ptr<arrow::Field>& PhysicalField() const override {
+ return physical_field_;
+ }
+
+ private:
+ std::shared_ptr<arrow::Field> logical_field_;
+ std::shared_ptr<arrow::Field> physical_field_;
+};
+
+class FullMapReadPlan : public MapShreddingColumnReadPlan {
public:
FullMapReadPlan(const std::shared_ptr<arrow::Field>& logical_field,
const std::shared_ptr<arrow::Field>& physical_read_field,
std::vector<std::pair<std::string, int32_t>>&&
selected_key_ids)
- : MapFieldReadPlan(logical_field, physical_read_field),
+ : MapShreddingColumnReadPlan(logical_field, physical_read_field),
selected_key_ids_(std::move(selected_key_ids)),
logical_map_type_(checked_pointer_cast<arrow::MapType>(logical_field->type()))
{}
- Result<std::shared_ptr<arrow::Array>> Materialize(
+ Result<std::shared_ptr<arrow::Array>> Assemble(
const std::shared_ptr<arrow::Array>& physical_array,
arrow::MemoryPool* arrow_pool) const override;
@@ -87,7 +101,7 @@ class FullMapReadPlan : public MapFieldReadPlan {
std::shared_ptr<arrow::MapType> logical_map_type_;
};
-class SharedSelectedKeysReadPlan : public MapFieldReadPlan {
+class SharedSelectedKeysReadPlan : public MapShreddingColumnReadPlan {
public:
struct SelectedKey {
int32_t field_id = -1;
@@ -98,10 +112,10 @@ class SharedSelectedKeysReadPlan : public MapFieldReadPlan
{
SharedSelectedKeysReadPlan(const std::shared_ptr<arrow::Field>&
logical_field,
const std::shared_ptr<arrow::Field>&
physical_read_field,
std::vector<SelectedKey>&& selected_keys)
- : MapFieldReadPlan(logical_field, physical_read_field),
+ : MapShreddingColumnReadPlan(logical_field, physical_read_field),
selected_keys_(std::move(selected_keys)) {}
- Result<std::shared_ptr<arrow::Array>> Materialize(
+ Result<std::shared_ptr<arrow::Array>> Assemble(
const std::shared_ptr<arrow::Array>& physical_array,
arrow::MemoryPool* arrow_pool) const override;
@@ -159,14 +173,15 @@ Result<std::shared_ptr<arrow::Array>>
MaskSinglePhysicalColumn(
return arrow::MakeArray(std::move(result_data));
}
-class DefaultSelectedKeysReadPlan : public MapFieldReadPlan {
+class DefaultSelectedKeysReadPlan : public MapShreddingColumnReadPlan {
public:
DefaultSelectedKeysReadPlan(const std::shared_ptr<arrow::Field>&
logical_field,
const std::shared_ptr<arrow::Field>&
physical_read_field,
const std::vector<std::string>& selected_keys)
- : MapFieldReadPlan(logical_field, physical_read_field),
selected_keys_(selected_keys) {}
+ : MapShreddingColumnReadPlan(logical_field, physical_read_field),
+ selected_keys_(selected_keys) {}
- Result<std::shared_ptr<arrow::Array>> Materialize(
+ Result<std::shared_ptr<arrow::Array>> Assemble(
const std::shared_ptr<arrow::Array>& physical_array,
arrow::MemoryPool* arrow_pool) const override;
@@ -176,7 +191,8 @@ class DefaultSelectedKeysReadPlan : public MapFieldReadPlan
{
} // namespace
-Result<std::unique_ptr<MapFieldReadPlan>>
MapFieldReadPlanFactory::CreateMapReadPlan(
+Result<std::shared_ptr<ShreddingColumnReadPlan>>
+MapSharedShreddingReadPlanFactory::CreateMapReadPlan(
const std::shared_ptr<arrow::Field>& logical_map_field,
const MapSharedShreddingFieldMeta& meta) {
if (logical_map_field->type()->id() != arrow::Type::MAP) {
@@ -213,12 +229,13 @@ Result<std::unique_ptr<MapFieldReadPlan>>
MapFieldReadPlanFactory::CreateMapRead
logical_map_type->item_type(), selected_physical_column_ids,
logical_map_type->item_field()->nullable(), include_overflow);
auto physical_read_field = logical_map_field->WithType(physical_type);
- std::unique_ptr<MapFieldReadPlan> read_plan =
std::make_unique<FullMapReadPlan>(
+ std::shared_ptr<ShreddingColumnReadPlan> read_plan =
std::make_shared<FullMapReadPlan>(
logical_map_field, physical_read_field, ResolveSelectedKeyIds(meta,
selected_keys));
return read_plan;
}
-Result<std::unique_ptr<MapFieldReadPlan>>
MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan(
+Result<std::shared_ptr<ShreddingColumnReadPlan>>
+MapSharedShreddingReadPlanFactory::CreateSharedSelectedKeysReadPlan(
const std::shared_ptr<arrow::Field>& selected_keys_field,
const MapSharedShreddingFieldMeta& meta) {
PAIMON_ASSIGN_OR_RAISE(
@@ -253,13 +270,14 @@ Result<std::unique_ptr<MapFieldReadPlan>>
MapFieldReadPlanFactory::CreateSharedS
value_field->type(), selected_physical_column_ids,
value_field->nullable(),
include_overflow);
auto physical_read_field = selected_keys_field->WithType(physical_type);
- std::unique_ptr<MapFieldReadPlan> read_plan =
std::make_unique<SharedSelectedKeysReadPlan>(
- selected_keys_field, physical_read_field,
std::move(selected_key_plans));
+ std::shared_ptr<ShreddingColumnReadPlan> read_plan =
+ std::make_shared<SharedSelectedKeysReadPlan>(selected_keys_field,
physical_read_field,
+
std::move(selected_key_plans));
return read_plan;
}
-Result<std::unique_ptr<MapFieldReadPlan>>
-MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
+Result<std::shared_ptr<ShreddingColumnReadPlan>>
+MapSharedShreddingReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
const std::shared_ptr<arrow::Field>& file_map_field,
const std::shared_ptr<arrow::Field>& selected_keys_field) {
if (file_map_field->type()->id() != arrow::Type::MAP) {
@@ -271,145 +289,13 @@
MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
std::vector<std::string> selected_keys,
NestedProjectionUtils::ValidateMapSharedShreddingAccessField(selected_keys_field));
auto physical_read_field =
selected_keys_field->WithType(file_map_field->type());
- std::unique_ptr<MapFieldReadPlan> read_plan =
std::make_unique<DefaultSelectedKeysReadPlan>(
- selected_keys_field, physical_read_field, selected_keys);
+ std::shared_ptr<ShreddingColumnReadPlan> read_plan =
+ std::make_shared<DefaultSelectedKeysReadPlan>(selected_keys_field,
physical_read_field,
+ selected_keys);
return read_plan;
}
-MapSharedShreddingFileReader::MapSharedShreddingFileReader(
- std::unique_ptr<FileBatchReader>&& reader,
- std::map<std::string, std::unique_ptr<MapFieldReadPlan>>&&
field_read_plans,
- const std::shared_ptr<MemoryPool>& pool)
- : arrow_pool_(GetArrowPool(pool)),
- reader_(std::move(reader)),
- field_read_plans_(std::move(field_read_plans)) {}
-
-Result<std::unique_ptr<::ArrowSchema>>
MapSharedShreddingFileReader::GetFileSchema() const {
- PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<::ArrowSchema> physical_schema,
- reader_->GetFileSchema());
- PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Schema>
physical_arrow_schema,
-
arrow::ImportSchema(physical_schema.get()));
-
- arrow::FieldVector logical_fields = physical_arrow_schema->fields();
- for (int32_t i = 0; i < physical_arrow_schema->num_fields(); ++i) {
- const auto& field = physical_arrow_schema->field(i);
- std::shared_ptr<arrow::KeyValueMetadata> metadata =
-
std::const_pointer_cast<arrow::KeyValueMetadata>(field->metadata());
- if (!MapSharedShreddingUtils::HasShreddingMetadata(metadata)) {
- continue;
- }
- PAIMON_ASSIGN_OR_RAISE(logical_fields[i], ToLogicalMapField(field));
- }
-
- auto logical_schema = arrow::schema(std::move(logical_fields));
- auto c_logical_schema = std::make_unique<ArrowSchema>();
- PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*logical_schema,
c_logical_schema.get()));
- return c_logical_schema;
-}
-
-Result<std::shared_ptr<arrow::Field>>
MapSharedShreddingFileReader::ToLogicalMapField(
- const std::shared_ptr<arrow::Field>& physical_field) {
- if (!physical_field || !physical_field->type() ||
- physical_field->type()->id() != arrow::Type::STRUCT) {
- return Status::Invalid(fmt::format("shared-shredding field {} is not a
physical struct",
- physical_field ?
physical_field->name() : "<null>"));
- }
- auto physical_type =
checked_pointer_cast<arrow::StructType>(physical_field->type());
- std::shared_ptr<arrow::DataType> value_type;
- bool value_nullable = true;
- for (const auto& child : physical_type->fields()) {
- if (child->name() == MapSharedShreddingDefine::kFieldMapping ||
- child->name() == MapSharedShreddingDefine::kOverflow) {
- continue;
- }
- value_type = child->type();
- value_nullable = child->nullable();
- break;
- }
- if (!value_type) {
- return Status::Invalid(fmt::format("cannot infer shared-shredding
value type for field {}",
- physical_field->name()));
- }
- return arrow::field(
- physical_field->name(),
- arrow::map(arrow::utf8(), arrow::field("value", value_type,
value_nullable)),
- physical_field->nullable());
-}
-
-Status MapSharedShreddingFileReader::SetReadSchema(
- ::ArrowSchema* read_schema, const std::shared_ptr<Predicate>& predicate,
- const std::optional<RoaringBitmap32>& selection_bitmap) {
- if (!read_schema) {
- return Status::Invalid(
- "invalid read schema in MapSharedShreddingFileReader, cannot be
null");
- }
- PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Schema>
logical_read_schema,
- arrow::ImportSchema(read_schema));
- bool converted = false;
- arrow::FieldVector physical_read_fields = logical_read_schema->fields();
- for (size_t i = 0; i < logical_read_schema->fields().size(); ++i) {
- const auto& field = logical_read_schema->field(i);
- auto plan_iter = field_read_plans_.find(field->name());
- if (plan_iter != field_read_plans_.end()) {
- physical_read_fields[i] = plan_iter->second->PhysicalReadField();
- converted = true;
- }
- }
- if (!converted) {
- return Status::Invalid("suppose not fall into
MapSharedShreddingFileReader");
- }
- auto physical_read_schema = arrow::schema(std::move(physical_read_fields));
- std::unique_ptr<ArrowSchema> c_physical_read_schema =
std::make_unique<ArrowSchema>();
- PAIMON_RETURN_NOT_OK_FROM_ARROW(
- arrow::ExportSchema(*physical_read_schema,
c_physical_read_schema.get()));
- return reader_->SetReadSchema(c_physical_read_schema.get(), predicate,
selection_bitmap);
-}
-
-Result<BatchReader::ReadBatch> MapSharedShreddingFileReader::NextBatch() {
- return Status::Invalid(
- "paimon inner reader MapSharedShreddingFileReader should use
NextBatchWithBitmap");
-}
-
-Result<BatchReader::ReadBatchWithBitmap>
MapSharedShreddingFileReader::NextBatchWithBitmap() {
- PAIMON_ASSIGN_OR_RAISE(BatchReader::ReadBatchWithBitmap batch_with_bitmap,
- reader_->NextBatchWithBitmap());
- if (BatchReader::IsEofBatch(batch_with_bitmap)) {
- return batch_with_bitmap;
- }
-
- auto& [batch, bitmap] = batch_with_bitmap;
- auto& [c_array, c_schema] = batch;
- PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Array>
arrow_array,
- arrow::ImportArray(c_array.get(),
c_schema.get()));
- if (!arrow_array || arrow_array->type_id() != arrow::Type::STRUCT) {
- return Status::Invalid("cannot cast batch to StructArray in
MapSharedShreddingFileReader");
- }
- auto struct_array = checked_pointer_cast<arrow::StructArray>(arrow_array);
-
- arrow::ArrayVector resolved_arrays = struct_array->fields();
- arrow::FieldVector resolved_fields = struct_array->struct_type()->fields();
- for (int32_t field_idx = 0; field_idx < struct_array->num_fields();
++field_idx) {
- const auto& physical_field =
struct_array->struct_type()->field(field_idx);
- auto plan_iter = field_read_plans_.find(physical_field->name());
- if (plan_iter == field_read_plans_.end()) {
- continue;
- }
- PAIMON_ASSIGN_OR_RAISE(
- resolved_arrays[field_idx],
- plan_iter->second->Materialize(struct_array->field(field_idx),
arrow_pool_.get()));
- resolved_fields[field_idx] = plan_iter->second->LogicalField();
- }
- PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::StructArray>
new_struct_array,
-
arrow::StructArray::Make(resolved_arrays, resolved_fields));
- auto new_c_array = std::make_unique<ArrowArray>();
- auto new_c_schema = std::make_unique<ArrowSchema>();
- PAIMON_RETURN_NOT_OK_FROM_ARROW(
- arrow::ExportArray(*new_struct_array, new_c_array.get(),
new_c_schema.get()));
- batch = std::make_pair(std::move(new_c_array), std::move(new_c_schema));
- return batch_with_bitmap;
-}
-
-Result<std::shared_ptr<arrow::Array>> FullMapReadPlan::Materialize(
+Result<std::shared_ptr<arrow::Array>> FullMapReadPlan::Assemble(
const std::shared_ptr<arrow::Array>& physical_array, arrow::MemoryPool*
arrow_pool) const {
if (!physical_array || physical_array->type_id() != arrow::Type::STRUCT) {
return Status::Invalid(fmt::format("cannot cast physical shredding
field {} to StructArray",
@@ -547,7 +433,7 @@ Result<std::shared_ptr<arrow::Array>>
FullMapReadPlan::Materialize(
return map_array;
}
-Result<std::shared_ptr<arrow::Array>> SharedSelectedKeysReadPlan::Materialize(
+Result<std::shared_ptr<arrow::Array>> SharedSelectedKeysReadPlan::Assemble(
const std::shared_ptr<arrow::Array>& physical_array, arrow::MemoryPool*
arrow_pool) const {
if (!physical_array || physical_array->type_id() != arrow::Type::STRUCT) {
return Status::Invalid(fmt::format("cannot cast physical shredding
field {} to StructArray",
@@ -716,7 +602,7 @@ Result<std::shared_ptr<arrow::Array>>
SharedSelectedKeysReadPlan::Materialize(
return result;
}
-Result<std::shared_ptr<arrow::Array>> DefaultSelectedKeysReadPlan::Materialize(
+Result<std::shared_ptr<arrow::Array>> DefaultSelectedKeysReadPlan::Assemble(
const std::shared_ptr<arrow::Array>& physical_array, arrow::MemoryPool*
arrow_pool) const {
if (!physical_array || physical_array->type_id() != arrow::Type::MAP) {
return Status::Invalid(
@@ -774,25 +660,4 @@ Result<std::shared_ptr<arrow::Array>>
DefaultSelectedKeysReadPlan::Materialize(
return result;
}
-std::shared_ptr<Metrics> MapSharedShreddingFileReader::GetReaderMetrics()
const {
- return reader_->GetReaderMetrics();
-}
-
-void MapSharedShreddingFileReader::Close() {
- reader_->Close();
-}
-
-Result<uint64_t> MapSharedShreddingFileReader::GetPreviousBatchFileRowId(
- uint64_t batch_row_id) const {
- return reader_->GetPreviousBatchFileRowId(batch_row_id);
-}
-
-Result<uint64_t> MapSharedShreddingFileReader::GetNumberOfRows() const {
- return reader_->GetNumberOfRows();
-}
-
-bool MapSharedShreddingFileReader::SupportPreciseBitmapSelection() const {
- return reader_->SupportPreciseBitmapSelection();
-}
-
} // namespace paimon
diff --git
a/src/paimon/common/data/shredding/map_shared_shredding_read_plan_factory.h
b/src/paimon/common/data/shredding/map_shared_shredding_read_plan_factory.h
new file mode 100644
index 00000000..871c51c3
--- /dev/null
+++ b/src/paimon/common/data/shredding/map_shared_shredding_read_plan_factory.h
@@ -0,0 +1,49 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+#pragma once
+
+#include <memory>
+
+#include "arrow/api.h"
+#include "paimon/common/data/shredding/map_shared_shredding_utils.h"
+#include "paimon/common/data/shredding/shredding_read_plan.h"
+
+namespace paimon {
+
+/// Builds per-column read plans for shared-shredding MAP columns and
selected-key MAP access.
+class MapSharedShreddingReadPlanFactory {
+ public:
+ MapSharedShreddingReadPlanFactory() = delete;
+ ~MapSharedShreddingReadPlanFactory() = delete;
+
+ static Result<std::shared_ptr<ShreddingColumnReadPlan>> CreateMapReadPlan(
+ const std::shared_ptr<arrow::Field>& logical_map_field,
+ const MapSharedShreddingFieldMeta& meta);
+
+ static Result<std::shared_ptr<ShreddingColumnReadPlan>>
CreateSharedSelectedKeysReadPlan(
+ const std::shared_ptr<arrow::Field>& selected_keys_field,
+ const MapSharedShreddingFieldMeta& meta);
+
+ static Result<std::shared_ptr<ShreddingColumnReadPlan>>
CreateDefaultSelectedKeysReadPlan(
+ const std::shared_ptr<arrow::Field>& file_map_field,
+ const std::shared_ptr<arrow::Field>& selected_keys_field);
+};
+
+} // namespace paimon
diff --git
a/src/paimon/common/data/shredding/map_shared_shredding_file_reader_test.cpp
b/src/paimon/common/data/shredding/map_shared_shredding_read_plan_factory_test.cpp
similarity index 90%
rename from
src/paimon/common/data/shredding/map_shared_shredding_file_reader_test.cpp
rename to
src/paimon/common/data/shredding/map_shared_shredding_read_plan_factory_test.cpp
index b2dd9374..e02f0ffb 100644
--- a/src/paimon/common/data/shredding/map_shared_shredding_file_reader_test.cpp
+++
b/src/paimon/common/data/shredding/map_shared_shredding_read_plan_factory_test.cpp
@@ -17,7 +17,7 @@
* under the License.
*/
-#include "paimon/common/data/shredding/map_shared_shredding_file_reader.h"
+#include
"paimon/common/data/shredding/map_shared_shredding_read_plan_factory.h"
#include <map>
#include <memory>
@@ -33,6 +33,7 @@
#include "gtest/gtest.h"
#include "paimon/common/data/shredding/map_shared_shredding_utils.h"
#include "paimon/common/data/shredding/map_shredding_defs.h"
+#include "paimon/common/data/shredding/shredding_file_reader.h"
#include "paimon/common/fs/external_path_provider.h"
#include "paimon/common/utils/checked_cast.h"
#include "paimon/core/append/append_only_writer.h"
@@ -51,7 +52,7 @@
#include "paimon/testing/utils/testharness.h"
namespace paimon::test {
-class MapSharedShreddingFileReaderTest : public ::testing::Test {
+class MapSharedShreddingReadPlanFactoryTest : public ::testing::Test {
public:
void SetUp() override {
pool_ = GetDefaultPool();
@@ -100,12 +101,12 @@ class MapSharedShreddingFileReaderTest : public
::testing::Test {
.ValueOrDie();
}
- std::unique_ptr<MapSharedShreddingFileReader> WrapReader(
+ std::unique_ptr<ShreddingFileReader> WrapReader(
std::unique_ptr<FileBatchReader>&& reader,
const std::optional<std::string>& selected_keys_str = std::nullopt)
const {
EXPECT_OK_AND_ASSIGN(auto c_file_schema, reader->GetFileSchema());
auto file_schema =
arrow::ImportSchema(c_file_schema.get()).ValueOrDie();
- std::map<std::string, std::unique_ptr<MapFieldReadPlan>>
field_read_plans;
+ std::map<std::string, std::shared_ptr<ShreddingColumnReadPlan>>
field_read_plans;
for (const auto& field : file_schema->fields()) {
auto metadata =
std::const_pointer_cast<arrow::KeyValueMetadata>(field->metadata());
if (!MapSharedShreddingUtils::HasShreddingMetadata(metadata)) {
@@ -124,20 +125,22 @@ class MapSharedShreddingFileReaderTest : public
::testing::Test {
EXPECT_TRUE(item_field);
auto map_type = checked_pointer_cast<arrow::MapType>(arrow::map(
arrow::utf8(), arrow::field("value", item_field->type(),
item_field->nullable())));
- std::shared_ptr<arrow::Field> logical_map_field =
field->WithType(map_type);
+ std::shared_ptr<arrow::Field> logical_map_field =
+ arrow::field(field->name(), map_type, field->nullable());
if (selected_keys_str.has_value()) {
logical_map_field =
logical_map_field->WithMetadata(arrow::KeyValueMetadata::Make(
{DataField::MAP_SELECTED_KEYS},
{selected_keys_str.value()}));
}
- EXPECT_OK_AND_ASSIGN(auto field_read_plan,
MapFieldReadPlanFactory::CreateMapReadPlan(
- logical_map_field,
meta));
+ EXPECT_OK_AND_ASSIGN(
+ auto field_read_plan,
+
MapSharedShreddingReadPlanFactory::CreateMapReadPlan(logical_map_field, meta));
field_read_plans.emplace(field->name(),
std::move(field_read_plan));
}
- return
std::make_unique<MapSharedShreddingFileReader>(std::move(reader),
-
std::move(field_read_plans), pool_);
+ return std::make_unique<ShreddingFileReader>(std::move(reader),
std::move(field_read_plans),
+ pool_);
}
- Result<std::unique_ptr<MapSharedShreddingFileReader>> CreateReader(
+ Result<std::unique_ptr<ShreddingFileReader>> CreateReader(
std::shared_ptr<arrow::Array> physical_array = nullptr,
std::shared_ptr<arrow::Schema> physical_schema = nullptr,
const std::optional<std::string>& selected_keys = std::nullopt) const {
@@ -236,20 +239,7 @@ class MapSharedShreddingFileReaderTest : public
::testing::Test {
};
};
-TEST_F(MapSharedShreddingFileReaderTest,
TestGetFileSchemaReturnsLogicalMapSchema) {
- ASSERT_OK_AND_ASSIGN(auto reader, CreateReader());
-
- ASSERT_OK_AND_ASSIGN(auto c_schema, reader->GetFileSchema());
- auto schema = arrow::ImportSchema(c_schema.get()).ValueOrDie();
-
- ASSERT_TRUE(schema->Equals(logical_schema_, /*check_metadata=*/false))
- << "Expected:\n"
- << logical_schema_->ToString() << "\nActual:\n"
- << schema->ToString();
- ASSERT_FALSE(schema->field(1)->HasMetadata());
-}
-
-TEST_F(MapSharedShreddingFileReaderTest,
TestAllExistSelectedKeysWithoutOverflow) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestAllExistSelectedKeysWithoutOverflow) {
ASSERT_OK_AND_ASSIGN(auto reader,
CreateReader(/*physical_array=*/nullptr,
/*physical_schema=*/nullptr,
/*selected_keys=*/"b"));
@@ -271,7 +261,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestAllExistSelectedKeysWithoutOverflow
AssertChunkedArrayEquals(expected, actual);
}
-TEST_F(MapSharedShreddingFileReaderTest, TestAllExistSelectedKeysWithOverflow)
{
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestAllExistSelectedKeysWithOverflow) {
ASSERT_OK_AND_ASSIGN(auto reader,
CreateReader(/*physical_array=*/nullptr,
/*physical_schema=*/nullptr,
/*selected_keys=*/"a,c"));
@@ -293,7 +283,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestAllExistSelectedKeysWithOverflow) {
AssertChunkedArrayEquals(expected, actual);
}
-TEST_F(MapSharedShreddingFileReaderTest, TestSelectedKeysStructProjection) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestSelectedKeysStructProjection) {
ASSERT_OK_AND_ASSIGN(auto physical_schema, PhysicalSchemaWithMetadata());
ASSERT_OK_AND_ASSIGN(auto physical_array, PhysicalArray());
auto mock_reader = std::make_unique<MockFileBatchReader>(
@@ -306,13 +296,13 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjection) {
auto selected_field = arrow::field(
"tags", selected_type, /*nullable=*/true,
arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS},
{"a,c,missing"}));
- ASSERT_OK_AND_ASSIGN(
- auto field_read_plan,
-
MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan(selected_field,
TagsMeta()));
- std::map<std::string, std::unique_ptr<MapFieldReadPlan>> contexts;
+ ASSERT_OK_AND_ASSIGN(auto field_read_plan,
+
MapSharedShreddingReadPlanFactory::CreateSharedSelectedKeysReadPlan(
+ selected_field, TagsMeta()));
+ std::map<std::string, std::shared_ptr<ShreddingColumnReadPlan>> contexts;
contexts.emplace("tags", std::move(field_read_plan));
- auto reader =
std::make_unique<MapSharedShreddingFileReader>(std::move(mock_reader),
-
std::move(contexts), pool_);
+ auto reader =
+ std::make_unique<ShreddingFileReader>(std::move(mock_reader),
std::move(contexts), pool_);
auto read_schema =
ExportSchema(arrow::schema({arrow::field("id", arrow::int32()),
selected_field}));
@@ -333,7 +323,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjection) {
AssertChunkedArrayEquals(expected, actual);
}
-TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjectionSharesValueBuffers) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestSelectedKeysStructProjectionSharesValueBuffers) {
ASSERT_OK_AND_ASSIGN(auto physical_array, PhysicalArray());
auto physical_root =
checked_pointer_cast<arrow::StructArray>(physical_array);
auto physical_tags =
@@ -347,11 +337,11 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjectionSharesV
auto selected_field = arrow::field(
"tags", selected_type, /*nullable=*/true,
arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS},
{"a,b,e,missing"}));
- ASSERT_OK_AND_ASSIGN(
- auto field_read_plan,
-
MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan(selected_field,
TagsMeta()));
+ ASSERT_OK_AND_ASSIGN(auto field_read_plan,
+
MapSharedShreddingReadPlanFactory::CreateSharedSelectedKeysReadPlan(
+ selected_field, TagsMeta()));
ASSERT_OK_AND_ASSIGN(auto result,
- field_read_plan->Materialize(physical_tags,
arrow::default_memory_pool()));
+ field_read_plan->Assemble(physical_tags,
arrow::default_memory_pool()));
auto result_struct = checked_pointer_cast<arrow::StructArray>(result);
auto expected = arrow::ipc::internal::json::ArrayFromJSON(selected_type,
R"([
@@ -370,7 +360,8 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjectionSharesV
ASSERT_EQ(physical_tags->data()->buffers[0],
result_struct->data()->buffers[0]);
}
-TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjectionSharesNestedValueBuffers) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
+ TestSelectedKeysStructProjectionSharesNestedValueBuffers) {
auto item_type = arrow::list(arrow::int64());
auto logical_schema =
arrow::schema({arrow::field("id", arrow::int32()),
@@ -402,9 +393,9 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjectionSharesN
arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"b"}));
ASSERT_OK_AND_ASSIGN(
auto field_read_plan,
-
MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan(selected_field,
meta));
+
MapSharedShreddingReadPlanFactory::CreateSharedSelectedKeysReadPlan(selected_field,
meta));
ASSERT_OK_AND_ASSIGN(auto result,
- field_read_plan->Materialize(physical_tags,
arrow::default_memory_pool()));
+ field_read_plan->Assemble(physical_tags,
arrow::default_memory_pool()));
auto result_struct = checked_pointer_cast<arrow::StructArray>(result);
auto result_list =
checked_pointer_cast<arrow::ListArray>(result_struct->field(0));
@@ -423,7 +414,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjectionSharesN
result_list->data()->child_data[0]->buffers[1]);
}
-TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjectionFromDefaultMap) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestSelectedKeysStructProjectionFromDefaultMap) {
auto map_type = checked_pointer_cast<arrow::MapType>(
arrow::map(arrow::utf8(), arrow::field("value", arrow::int64())));
auto file_schema =
@@ -446,12 +437,12 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjectionFromDef
arrow::field("tags", selected_type, /*nullable=*/true,
arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a,missing"}));
ASSERT_OK_AND_ASSIGN(auto field_read_plan,
-
MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
+
MapSharedShreddingReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
file_schema->field(1), selected_field));
- std::map<std::string, std::unique_ptr<MapFieldReadPlan>> contexts;
+ std::map<std::string, std::shared_ptr<ShreddingColumnReadPlan>> contexts;
contexts.emplace("tags", std::move(field_read_plan));
- auto reader =
std::make_unique<MapSharedShreddingFileReader>(std::move(mock_reader),
-
std::move(contexts), pool_);
+ auto reader =
+ std::make_unique<ShreddingFileReader>(std::move(mock_reader),
std::move(contexts), pool_);
auto read_schema =
ExportSchema(arrow::schema({arrow::field("id", arrow::int32()),
selected_field}));
@@ -471,15 +462,15 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestSelectedKeysStructProjectionFromDef
AssertChunkedArrayEquals(expected, actual);
}
-TEST_F(MapSharedShreddingFileReaderTest,
TestInvalidSelectedKeysStructProjection) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestInvalidSelectedKeysStructProjection) {
auto file_map_field = arrow::field("tags", arrow::map(arrow::utf8(),
arrow::int64()));
auto mismatched_count_field =
arrow::field("tags", arrow::struct_({arrow::field("a",
arrow::int64())}), /*nullable=*/true,
arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a,b"}));
-
ASSERT_NOK_WITH_MSG(MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan(
+
ASSERT_NOK_WITH_MSG(MapSharedShreddingReadPlanFactory::CreateSharedSelectedKeysReadPlan(
mismatched_count_field, TagsMeta()),
"metadata size 2 does not match STRUCT field count 1");
-
ASSERT_NOK_WITH_MSG(MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
+
ASSERT_NOK_WITH_MSG(MapSharedShreddingReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
file_map_field, mismatched_count_field),
"metadata size 2 does not match STRUCT field count 1");
@@ -487,15 +478,15 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestInvalidSelectedKeysStructProjection
"tags",
arrow::struct_({arrow::field("a", arrow::int64()), arrow::field("b",
arrow::utf8())}),
/*nullable=*/true,
arrow::KeyValueMetadata::Make({DataField::MAP_SELECTED_KEYS}, {"a,b"}));
-
ASSERT_NOK_WITH_MSG(MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan(
+
ASSERT_NOK_WITH_MSG(MapSharedShreddingReadPlanFactory::CreateSharedSelectedKeysReadPlan(
mismatched_type_field, TagsMeta()),
"must have the same value type");
-
ASSERT_NOK_WITH_MSG(MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
+
ASSERT_NOK_WITH_MSG(MapSharedShreddingReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
file_map_field, mismatched_type_field),
"must have the same value type");
}
-TEST_F(MapSharedShreddingFileReaderTest, TestPartialExistSelectedKeys) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest, TestPartialExistSelectedKeys) {
ASSERT_OK_AND_ASSIGN(auto reader,
CreateReader(/*physical_array=*/nullptr,
/*physical_schema=*/nullptr,
/*selected_keys=*/"a,c,missing"));
@@ -518,7 +509,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestPartialExistSelectedKeys) {
AssertChunkedArrayEquals(expected, actual);
}
-TEST_F(MapSharedShreddingFileReaderTest, TestMissingSelectedKeysReadsWholeMap)
{
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestMissingSelectedKeysReadsWholeMap) {
ASSERT_OK_AND_ASSIGN(auto reader, CreateReader());
auto read_schema = ExportSchema(ReadSchema(std::nullopt));
ASSERT_OK(reader->SetReadSchema(read_schema.get(), /*predicate=*/nullptr,
@@ -538,7 +529,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestMissingSelectedKeysReadsWholeMap) {
AssertChunkedArrayEquals(expected, actual);
}
-TEST_F(MapSharedShreddingFileReaderTest, TestSpecialSelectedKeys) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest, TestSpecialSelectedKeys) {
MapSharedShreddingFieldMeta meta;
meta.name_to_id = {{"", 0}, {" ", 1}, {".", 2}, {"a", 3}};
meta.field_to_columns = {{0, {0}}, {1, {1}}, {2, {0}}, {3, {1}}};
@@ -591,7 +582,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestSpecialSelectedKeys) {
])");
}
-TEST_F(MapSharedShreddingFileReaderTest,
TestUnknownSelectedKeyReturnsEmptyMap) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestUnknownSelectedKeyReturnsEmptyMap) {
ASSERT_OK_AND_ASSIGN(auto reader,
CreateReader(/*physical_array=*/nullptr,
/*physical_schema=*/nullptr,
/*selected_keys=*/"missing"));
@@ -614,7 +605,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestUnknownSelectedKeyReturnsEmptyMap)
AssertChunkedArrayEquals(expected, actual);
}
-TEST_F(MapSharedShreddingFileReaderTest, TestInvalidNullFieldMappingField) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestInvalidNullFieldMappingField) {
ASSERT_OK_AND_ASSIGN(auto physical_schema, PhysicalSchemaWithMetadata());
std::string json = R"([
[1, [null, 10, null, null]]
@@ -631,7 +622,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestInvalidNullFieldMappingField) {
"__field_mapping cannot be null");
}
-TEST_F(MapSharedShreddingFileReaderTest,
TestInvalidNullFieldMappingFieldElement) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestInvalidNullFieldMappingFieldElement) {
ASSERT_OK_AND_ASSIGN(auto physical_schema, PhysicalSchemaWithMetadata());
std::string json = R"([
[1, [[0, null], 10, null, null]]
@@ -648,7 +639,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestInvalidNullFieldMappingFieldElement
"__field_mapping element cannot be null");
}
-TEST_F(MapSharedShreddingFileReaderTest, TestListValue) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest, TestListValue) {
std::shared_ptr<arrow::Schema> logical_schema = arrow::schema({
arrow::field("id", arrow::int32()),
arrow::field("tags", arrow::map(arrow::utf8(),
arrow::list(arrow::int32()))),
@@ -704,7 +695,7 @@ TEST_F(MapSharedShreddingFileReaderTest, TestListValue) {
AssertChunkedArrayEquals(expected, actual);
}
-TEST_F(MapSharedShreddingFileReaderTest, TestOrcDictionaryEncodedStringValue) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestOrcDictionaryEncodedStringValue) {
std::shared_ptr<arrow::Schema> logical_schema = arrow::schema({
arrow::field("id", arrow::int32()),
arrow::field("tags", arrow::map(arrow::utf8(), arrow::utf8())),
@@ -764,7 +755,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestOrcDictionaryEncodedStringValue) {
AssertChunkedArrayEquals(expected, actual);
}
-TEST_F(MapSharedShreddingFileReaderTest,
TestOrcDictionaryEncodedStringListValue) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest,
TestOrcDictionaryEncodedStringListValue) {
std::shared_ptr<arrow::Schema> logical_schema = arrow::schema({
arrow::field("id", arrow::int32()),
arrow::field("tags", arrow::map(arrow::utf8(),
arrow::list(arrow::utf8()))),
@@ -824,7 +815,7 @@ TEST_F(MapSharedShreddingFileReaderTest,
TestOrcDictionaryEncodedStringListValue
AssertChunkedArrayEquals(expected, actual);
}
-TEST_F(MapSharedShreddingFileReaderTest, TestReadsRealFormatFile) {
+TEST_F(MapSharedShreddingReadPlanFactoryTest, TestReadsRealFormatFile) {
// TODO(lisizhuo.lsz): support other format
auto options = options_;
std::string format = "orc";
diff --git a/src/paimon/common/data/shredding/shredding_file_reader.cpp
b/src/paimon/common/data/shredding/shredding_file_reader.cpp
index 0f47350d..1c2e34f5 100644
--- a/src/paimon/common/data/shredding/shredding_file_reader.cpp
+++ b/src/paimon/common/data/shredding/shredding_file_reader.cpp
@@ -84,7 +84,7 @@ Result<BatchReader::ReadBatchWithBitmap>
ShreddingFileReader::NextBatchWithBitma
auto& [c_array, c_schema] = batch;
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Array>
arrow_array,
arrow::ImportArray(c_array.get(),
c_schema.get()));
- if (arrow_array->type_id() != arrow::Type::STRUCT) {
+ if (!arrow_array || arrow_array->type_id() != arrow::Type::STRUCT) {
return Status::Invalid("cannot cast batch to StructArray in
ShreddingFileReader");
}
auto struct_array = checked_pointer_cast<arrow::StructArray>(arrow_array);
diff --git a/src/paimon/core/operation/abstract_split_read.cpp
b/src/paimon/core/operation/abstract_split_read.cpp
index bb82c5d8..d057e1d7 100644
--- a/src/paimon/core/operation/abstract_split_read.cpp
+++ b/src/paimon/core/operation/abstract_split_read.cpp
@@ -28,11 +28,10 @@
#include "fmt/format.h"
#include "paimon/common/data/blob_defs.h"
#include "paimon/common/data/blob_utils.h"
-#include "paimon/common/data/shredding/map_shared_shredding_file_reader.h"
+#include
"paimon/common/data/shredding/map_shared_shredding_read_plan_factory.h"
#include "paimon/common/data/shredding/map_shared_shredding_utils.h"
#include "paimon/common/data/shredding/shredding_file_reader.h"
#include "paimon/common/data/variant/variant_shredding_read_plan_factory.h"
-#include "paimon/common/data/variant/variant_type_utils.h"
#include "paimon/common/reader/delegating_prefetch_reader.h"
#include "paimon/common/reader/predicate_batch_reader.h"
#include "paimon/common/reader/prefetch_file_batch_reader_impl.h"
@@ -221,13 +220,11 @@ Result<std::unique_ptr<FileBatchReader>>
AbstractSplitRead::CreateFieldMappingRe
}
std::set<int32_t> skip_map_selected_keys_filter_field_ids;
if (file_format_identifier != "blob") {
- std::pair<std::unique_ptr<FileBatchReader>, std::set<int32_t>>
shared_shredding_result;
- PAIMON_ASSIGN_OR_RAISE(shared_shredding_result,
ApplySharedShreddingReaderIfNeeded(
-
std::move(file_reader), read_schema));
- file_reader = std::move(shared_shredding_result.first);
- skip_map_selected_keys_filter_field_ids =
std::move(shared_shredding_result.second);
- PAIMON_ASSIGN_OR_RAISE(
- file_reader,
ApplyVariantShreddingReaderIfNeeded(std::move(file_reader), read_schema));
+ std::pair<std::unique_ptr<FileBatchReader>, std::set<int32_t>>
shredding_result;
+ PAIMON_ASSIGN_OR_RAISE(shredding_result,
+
ApplyShreddingReaderIfNeeded(std::move(file_reader), read_schema));
+ file_reader = std::move(shredding_result.first);
+ skip_map_selected_keys_filter_field_ids =
std::move(shredding_result.second);
}
if (NeedCompleteRowTrackingFields(options_.RowTrackingEnabled(),
read_schema)) {
// A blob file has no self-describing schema: its physical fields are
declared by the
@@ -260,15 +257,16 @@ Result<std::unique_ptr<FileBatchReader>>
AbstractSplitRead::CreateFieldMappingRe
}
Result<std::pair<std::unique_ptr<FileBatchReader>, std::set<int32_t>>>
-AbstractSplitRead::ApplySharedShreddingReaderIfNeeded(
+AbstractSplitRead::ApplyShreddingReaderIfNeeded(
std::unique_ptr<FileBatchReader>&& file_reader,
const std::shared_ptr<arrow::Schema>& read_schema) const {
- std::set<int32_t> handled_shared_shredding_field_ids;
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<::ArrowSchema> file_schema,
file_reader->GetFileSchema());
PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Schema>
file_arrow_schema,
arrow::ImportSchema(file_schema.get()));
- std::map<std::string, std::unique_ptr<MapFieldReadPlan>> field_read_plans;
+
+ std::set<int32_t> handled_shared_shredding_field_ids;
+ std::map<std::string, std::shared_ptr<ShreddingColumnReadPlan>> plans;
for (const auto& read_field : read_schema->fields()) {
const auto& field_name = read_field->name();
auto file_field = file_arrow_schema->GetFieldByName(field_name);
@@ -286,64 +284,46 @@ AbstractSplitRead::ApplySharedShreddingReaderIfNeeded(
continue;
}
- std::unique_ptr<MapFieldReadPlan> field_read_plan;
+ std::shared_ptr<ShreddingColumnReadPlan> plan;
if (is_shared_shredding_map_access) {
if (is_shared_shredding_file) {
PAIMON_ASSIGN_OR_RAISE(MapSharedShreddingFieldMeta meta,
MapSharedShreddingUtils::DeserializeMetadata(metadata));
PAIMON_ASSIGN_OR_RAISE(
- field_read_plan,
-
MapFieldReadPlanFactory::CreateSharedSelectedKeysReadPlan(read_field, meta));
+ plan,
MapSharedShreddingReadPlanFactory::CreateSharedSelectedKeysReadPlan(
+ read_field, meta));
} else {
- PAIMON_ASSIGN_OR_RAISE(field_read_plan,
-
MapFieldReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
- file_field, read_field));
+ PAIMON_ASSIGN_OR_RAISE(
+ plan,
MapSharedShreddingReadPlanFactory::CreateDefaultSelectedKeysReadPlan(
+ file_field, read_field));
}
} else {
PAIMON_ASSIGN_OR_RAISE(MapSharedShreddingFieldMeta meta,
MapSharedShreddingUtils::DeserializeMetadata(metadata));
- PAIMON_ASSIGN_OR_RAISE(field_read_plan,
-
MapFieldReadPlanFactory::CreateMapReadPlan(read_field, meta));
+ PAIMON_ASSIGN_OR_RAISE(
+ plan,
MapSharedShreddingReadPlanFactory::CreateMapReadPlan(read_field, meta));
}
- field_read_plans.emplace(field_name, std::move(field_read_plan));
+ plans.emplace(field_name, std::move(plan));
PAIMON_ASSIGN_OR_RAISE(int32_t field_id,
NestedProjectionUtils::GetPaimonFieldId(read_field));
handled_shared_shredding_field_ids.insert(field_id);
}
- if (!field_read_plans.empty()) {
- file_reader = std::make_unique<MapSharedShreddingFileReader>(
- std::move(file_reader), std::move(field_read_plans), pool_);
- }
- return std::make_pair(std::move(file_reader),
std::move(handled_shared_shredding_field_ids));
-}
-Result<std::unique_ptr<FileBatchReader>>
AbstractSplitRead::ApplyVariantShreddingReaderIfNeeded(
- std::unique_ptr<FileBatchReader>&& file_reader,
- const std::shared_ptr<arrow::Schema>& read_schema) const {
- bool has_variant_field = false;
- for (const auto& read_field : read_schema->fields()) {
- // Variant columns may be nested inside struct columns; a
variant-access projection also
- // matches because it carries the variant extension marker itself.
- if (VariantTypeUtils::ContainsVariantField(read_field)) {
- has_variant_field = true;
- break;
+ std::map<std::string, std::shared_ptr<ShreddingColumnReadPlan>>
variant_plans;
+ PAIMON_ASSIGN_OR_RAISE(variant_plans,
VariantShreddingReadPlanFactory::CreateReadPlans(
+ read_schema, file_arrow_schema,
pool_));
+ for (auto& [field_name, plan] : variant_plans) {
+ if (!plans.emplace(field_name, std::move(plan)).second) {
+ return Status::Invalid(
+ fmt::format("multiple shredding read plans exist for field
{}", field_name));
}
}
- if (!has_variant_field) {
- return std::move(file_reader);
- }
- PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<::ArrowSchema> file_schema,
- file_reader->GetFileSchema());
- PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(std::shared_ptr<arrow::Schema>
file_arrow_schema,
- arrow::ImportSchema(file_schema.get()));
- std::map<std::string, std::shared_ptr<ShreddingColumnReadPlan>> plans;
- PAIMON_ASSIGN_OR_RAISE(plans,
VariantShreddingReadPlanFactory::CreateReadPlans(
- read_schema, file_arrow_schema, pool_));
+
if (!plans.empty()) {
file_reader =
std::make_unique<ShreddingFileReader>(std::move(file_reader),
std::move(plans), pool_);
}
- return std::move(file_reader);
+ return std::make_pair(std::move(file_reader),
std::move(handled_shared_shredding_field_ids));
}
Result<std::vector<DataField>>
AbstractSplitRead::ProjectFieldsForRowTrackingAndDataEvolution(
diff --git a/src/paimon/core/operation/abstract_split_read.h
b/src/paimon/core/operation/abstract_split_read.h
index 27349fec..a56b48fd 100644
--- a/src/paimon/core/operation/abstract_split_read.h
+++ b/src/paimon/core/operation/abstract_split_read.h
@@ -117,16 +117,12 @@ class AbstractSplitRead : public SplitRead {
const std::optional<std::vector<Range>>& row_ranges,
const std::shared_ptr<DataFilePathFactory>& data_file_path_factory)
const;
+ /// The returned field ID set contains MAP fields handled by shredding
read plans. It tells
+ /// FieldMappingReader to skip its generic selected-key filtering because
the plans have already
+ /// applied any requested key selection.
Result<std::pair<std::unique_ptr<FileBatchReader>, std::set<int32_t>>>
- ApplySharedShreddingReaderIfNeeded(std::unique_ptr<FileBatchReader>&&
file_reader,
- const std::shared_ptr<arrow::Schema>&
read_schema) const;
-
- /// Wraps the reader with a `ShreddingFileReader` when any read variant
column needs
- /// reassembly or path extraction; a plain read of an unshredded variant
column is passed
- /// through untouched.
- Result<std::unique_ptr<FileBatchReader>>
ApplyVariantShreddingReaderIfNeeded(
- std::unique_ptr<FileBatchReader>&& file_reader,
- const std::shared_ptr<arrow::Schema>& read_schema) const;
+ ApplyShreddingReaderIfNeeded(std::unique_ptr<FileBatchReader>&&
file_reader,
+ const std::shared_ptr<arrow::Schema>&
read_schema) const;
static bool NeedCompleteRowTrackingFields(bool row_tracking_enabled,
const
std::shared_ptr<arrow::Schema>& read_schema);
diff --git a/src/paimon/core/operation/data_evolution_split_read.h
b/src/paimon/core/operation/data_evolution_split_read.h
index b4f80adb..fad0e674 100644
--- a/src/paimon/core/operation/data_evolution_split_read.h
+++ b/src/paimon/core/operation/data_evolution_split_read.h
@@ -64,7 +64,6 @@ struct DeletionFile;
/// ->(ConcatBatchReader across blob files | BlobFallbackBatchReader across
blob sequence layers)
///
->FieldMappingReader->(ApplyDeletionVectorBatchReader)->(ApplyBitmapIndexBatchReader)
/// ->(CompleteRowTrackingFieldsBatchReader)->(ShreddingFileReader)
-/// ->(MapSharedShreddingFileReader)
///
->(VectorFileBatchReader)->(DelegatingPrefetchReader)->(PrefetchFileBatchReader)->FormatReader
///
///
diff --git a/src/paimon/core/operation/merge_file_split_read.h
b/src/paimon/core/operation/merge_file_split_read.h
index d4bfa727..5003cb55 100644
--- a/src/paimon/core/operation/merge_file_split_read.h
+++ b/src/paimon/core/operation/merge_file_split_read.h
@@ -74,7 +74,6 @@ class MergeFunctionWrapper;
/// files->KeyValueProjectionReader/AsyncKeyValueProjectionReader
///
->DropDeleteReader->SortMergeReader->ConcatKeyValueRecordReader->KeyValueDataFileRecordReader
///
->FieldMappingReader->(ApplyDeletionVectorBatchReader)->(ShreddingFileReader)
-/// ->(MapSharedShreddingFileReader)
/// ->(DelegatingPrefetchReader)->(PrefetchFileBatchReader)->FormatReader
class MergeFileSplitRead : public AbstractSplitRead {
public:
diff --git a/src/paimon/core/operation/raw_file_split_read.h
b/src/paimon/core/operation/raw_file_split_read.h
index 6a97b9b3..93eab550 100644
--- a/src/paimon/core/operation/raw_file_split_read.h
+++ b/src/paimon/core/operation/raw_file_split_read.h
@@ -54,7 +54,7 @@ struct DeletionFile;
/// splits)->CompleteRowKindBatchReader->(PredicateBatchReader)
/// ->ConcatBatchReader across
///
files->FieldMappingReader->(ApplyBitmapIndexBatchReader)->(CompleteRowTrackingFieldsBatchReader)
-///
->(ShreddingFileReader)->(MapSharedShreddingFileReader)->(VectorFileBatchReader)
+/// ->(ShreddingFileReader)->(VectorFileBatchReader)
/// ->(DelegatingPrefetchReader)->(PrefetchFileBatchReader)->FormatReader
class RawFileSplitRead : public AbstractSplitRead {