SteNicholas commented on code in PR #376: URL: https://github.com/apache/paimon-cpp/pull/376#discussion_r4111455787
########## src/paimon/format/vortex/vortex_file_batch_reader.cpp: ########## @@ -0,0 +1,559 @@ +/* + * 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. + */ + +#include "paimon/format/vortex/vortex_file_batch_reader.h" + +#include <algorithm> +#include <limits> +#include <string> +#include <utility> + +#include "arrow/api.h" +#include "arrow/c/bridge.h" +#include "fmt/format.h" +#include "paimon/common/metrics/metrics_impl.h" +#include "paimon/common/predicate/predicate_filter.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/math.h" +#include "paimon/format/vortex/vortex_ffi_util.h" +#include "paimon/fs/file_system.h" + +namespace paimon::vortex { + +namespace { + +// Vortex's Arrow export represents strings and binaries with the view types (StringView / +// BinaryView). paimon schemas use the standard string/binary types, so map the view types back, +// recursing through nested types. Types that contain no views are returned unchanged. +std::shared_ptr<arrow::DataType> NormalizeViewType(const std::shared_ptr<arrow::DataType>& type) { + switch (type->id()) { + case arrow::Type::STRING_VIEW: + return arrow::utf8(); + case arrow::Type::BINARY_VIEW: + return arrow::binary(); + case arrow::Type::STRUCT: { + arrow::FieldVector fields; + fields.reserve(type->num_fields()); + for (const std::shared_ptr<arrow::Field>& field : type->fields()) { + fields.push_back(field->WithType(NormalizeViewType(field->type()))); + } + return arrow::struct_(fields); + } + case arrow::Type::LIST: + return arrow::list(type->field(0)->WithType(NormalizeViewType(type->field(0)->type()))); + case arrow::Type::LARGE_LIST: + return arrow::large_list( + type->field(0)->WithType(NormalizeViewType(type->field(0)->type()))); + case arrow::Type::FIXED_SIZE_LIST: + return arrow::fixed_size_list( + type->field(0)->WithType(NormalizeViewType(type->field(0)->type())), + checked_cast<const arrow::FixedSizeListType&>(*type).list_size()); + default: + return type; + } +} + +std::shared_ptr<arrow::Schema> NormalizeViewSchema(const std::shared_ptr<arrow::Schema>& schema) { + arrow::FieldVector fields; + fields.reserve(schema->num_fields()); + for (const std::shared_ptr<arrow::Field>& field : schema->fields()) { + fields.push_back(field->WithType(NormalizeViewType(field->type()))); + } + return arrow::schema(fields, schema->metadata()); +} + +// Rebuilds `array` with every StringView/BinaryView column (which Vortex's Arrow export produces) +// converted to the standard string/binary type. Arrow 17 has no cast kernel between the view and +// standard types, so the leaves are rebuilt with builders and the parents reconstructed, recursing +// through structs and lists. Arrays with no view-typed data are returned unchanged. `array` must be +// offset-normalized (offset 0) so the rebuilt parents and children stay consistent. +Result<std::shared_ptr<arrow::Array>> NormalizeViewArray(const std::shared_ptr<arrow::Array>& array, + arrow::MemoryPool* pool) { + switch (array->type_id()) { + case arrow::Type::STRING_VIEW: { + const auto& view = checked_cast<const arrow::StringViewArray&>(*array); + arrow::StringBuilder builder(pool); + for (int64_t i = 0; i < view.length(); ++i) { + if (view.IsNull(i)) { + PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.AppendNull()); + } else { + PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.Append(view.GetView(i))); + } + } + std::shared_ptr<arrow::Array> out; + PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.Finish(&out)); + return out; + } + case arrow::Type::BINARY_VIEW: { + const auto& view = checked_cast<const arrow::BinaryViewArray&>(*array); + arrow::BinaryBuilder builder(pool); + for (int64_t i = 0; i < view.length(); ++i) { + if (view.IsNull(i)) { + PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.AppendNull()); + } else { + PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.Append(view.GetView(i))); + } + } + std::shared_ptr<arrow::Array> out; + PAIMON_RETURN_NOT_OK_FROM_ARROW(builder.Finish(&out)); + return out; + } + case arrow::Type::STRUCT: { + const auto& struct_array = checked_cast<const arrow::StructArray&>(*array); + arrow::ArrayVector children; + arrow::FieldVector fields; + children.reserve(struct_array.num_fields()); + fields.reserve(struct_array.num_fields()); + bool changed = false; + for (int32_t i = 0; i < struct_array.num_fields(); ++i) { + const std::shared_ptr<arrow::Array>& child = struct_array.field(i); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::Array> normalized_child, + NormalizeViewArray(child, pool)); + changed = changed || normalized_child.get() != child.get(); + fields.push_back(array->type()->field(i)->WithType(normalized_child->type())); + children.push_back(std::move(normalized_child)); + } + if (!changed) { + return array; + } + return std::make_shared<arrow::StructArray>( + arrow::struct_(fields), struct_array.length(), children, struct_array.null_bitmap(), + struct_array.null_count(), struct_array.offset()); + } + case arrow::Type::LIST: { + const auto& list = checked_cast<const arrow::ListArray&>(*array); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::Array> values, + NormalizeViewArray(list.values(), pool)); + if (values.get() == list.values().get()) { + return array; + } + return std::make_shared<arrow::ListArray>( + arrow::list(array->type()->field(0)->WithType(values->type())), list.length(), + list.value_offsets(), values, list.null_bitmap(), list.null_count(), list.offset()); + } + case arrow::Type::LARGE_LIST: { + const auto& list = checked_cast<const arrow::LargeListArray&>(*array); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::Array> values, + NormalizeViewArray(list.values(), pool)); + if (values.get() == list.values().get()) { + return array; + } + return std::make_shared<arrow::LargeListArray>( + arrow::large_list(array->type()->field(0)->WithType(values->type())), list.length(), + list.value_offsets(), values, list.null_bitmap(), list.null_count(), list.offset()); + } + case arrow::Type::FIXED_SIZE_LIST: { + const auto& list = checked_cast<const arrow::FixedSizeListArray&>(*array); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::Array> values, + NormalizeViewArray(list.values(), pool)); + if (values.get() == list.values().get()) { + return array; + } + const auto& fsl_type = checked_cast<const arrow::FixedSizeListType&>(*array->type()); + return std::make_shared<arrow::FixedSizeListArray>( + arrow::fixed_size_list(array->type()->field(0)->WithType(values->type()), + fsl_type.list_size()), + list.length(), values, list.null_bitmap(), list.null_count(), list.offset()); + } + default: + return array; + } +} + +// Recursively projects `array` onto `target_type` (the read schema's type for this level), +// selecting struct children by field name and pruning nested struct/list/fixed-size-list elements +// so the returned array's type equals `target_type`. Returns `array` unchanged when its type +// already matches. Leaf types are expected to match after view normalization; an unexpected leaf +// mismatch is returned as-is for the caller's schema check to surface. +Result<std::shared_ptr<arrow::Array>> ProjectArrayToType( + const std::shared_ptr<arrow::Array>& array, + const std::shared_ptr<arrow::DataType>& target_type) { + if (array->type()->Equals(target_type)) { + return array; + } + switch (target_type->id()) { + case arrow::Type::STRUCT: { + const auto& struct_array = checked_cast<const arrow::StructArray&>(*array); + const auto& struct_type = checked_cast<const arrow::StructType&>(*array->type()); + const auto& target_struct = checked_cast<const arrow::StructType&>(*target_type); + arrow::ArrayVector children; + children.reserve(target_struct.num_fields()); + for (const std::shared_ptr<arrow::Field>& target_field : target_struct.fields()) { + const int32_t index = struct_type.GetFieldIndex(target_field->name()); + if (index < 0) { + return Status::Invalid( + fmt::format("Vortex read field '{}' is not present in the file schema", + target_field->name())); + } + PAIMON_ASSIGN_OR_RAISE( + std::shared_ptr<arrow::Array> child, + ProjectArrayToType(struct_array.field(index), target_field->type())); + children.push_back(std::move(child)); + } + return std::make_shared<arrow::StructArray>( + target_type, struct_array.length(), children, struct_array.null_bitmap(), Review Comment: **Reject incompatible projected leaf types before exporting**. `ProjectArrayToType` returns the original array for a leaf-type mismatch, but this constructor assigns the requested parent type without converting its children. A direct `SetReadSchema` request for `int64` on an `int32` file column therefore exports 32-bit buffers as 64-bit values. With the extracted helper and Arrow 17, `[1,2,3]` round-tripped through the C interface as `[8589934593,3,0]`, while `ExportArray` reported success. There is no later schema check that rejects this. Please reject incompatible leaf/container types (or convert supported representations) before rebuilding the parent; the existing `NestedProjectionUtils::AlignArrayToReadType` already validates these cases. -- 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]
