diff --git a/be/src/format/arrow/arrow_row_batch.cpp b/be/src/format/arrow/arrow_row_batch.cpp index cc88bc18b92f29..9baf66dc927166 100644 --- a/be/src/format/arrow/arrow_row_batch.cpp +++ b/be/src/format/arrow/arrow_row_batch.cpp @@ -132,33 +132,47 @@ Status convert_to_arrow_type(const DataTypePtr& origin_type, break; case TYPE_ARRAY: { const auto* type_arr = assert_cast(remove_nullable(type).get()); + const auto& item_doris_type = type_arr->get_nested_type(); std::shared_ptr item_type; - RETURN_IF_ERROR(convert_to_arrow_type(type_arr->get_nested_type(), &item_type, timezone, - datetime_naive)); - *result = std::make_shared(item_type); + RETURN_IF_ERROR( + convert_to_arrow_type(item_doris_type, &item_type, timezone, datetime_naive)); + // Arrow keeps metadata on the Field, not on the DataType, so the element's Doris type has + // to be recorded here: ListType(item_type) would synthesize a bare "item" field and a + // nested LARGEINT would be indistinguishable from a nested STRING. The name and the + // nullability are the ones that constructor would have used, only the metadata is new. + *result = std::make_shared(create_arrow_field_with_metadata( + "item", item_type, /*is_nullable=*/true, item_doris_type->get_primitive_type())); break; } case TYPE_MAP: { const auto* type_map = assert_cast(remove_nullable(type).get()); + const auto& key_doris_type = type_map->get_key_type(); + const auto& val_doris_type = type_map->get_value_type(); std::shared_ptr key_type; std::shared_ptr val_type; - RETURN_IF_ERROR(convert_to_arrow_type(type_map->get_key_type(), &key_type, timezone, - datetime_naive)); - RETURN_IF_ERROR(convert_to_arrow_type(type_map->get_value_type(), &val_type, timezone, - datetime_naive)); - *result = std::make_shared(key_type, val_type); + RETURN_IF_ERROR(convert_to_arrow_type(key_doris_type, &key_type, timezone, datetime_naive)); + RETURN_IF_ERROR(convert_to_arrow_type(val_doris_type, &val_type, timezone, datetime_naive)); + // Same reason as the list element above. An Arrow map key is never nullable -- that is what + // MapType(key_type, val_type) builds and what the IPC format allows -- so only the metadata + // changes here as well. + *result = std::make_shared( + create_arrow_field_with_metadata("key", key_type, /*is_nullable=*/false, + key_doris_type->get_primitive_type()), + create_arrow_field_with_metadata("value", val_type, /*is_nullable=*/true, + val_doris_type->get_primitive_type())); break; } case TYPE_STRUCT: { const auto* type_struct = assert_cast(remove_nullable(type).get()); std::vector> fields; for (size_t i = 0; i < type_struct->get_elements().size(); i++) { + const auto& element_doris_type = type_struct->get_element(i); std::shared_ptr field_type; - RETURN_IF_ERROR(convert_to_arrow_type(type_struct->get_element(i), &field_type, - timezone, datetime_naive)); - fields.push_back( - std::make_shared(type_struct->get_element_name(i), field_type, - type_struct->get_element(i)->is_nullable())); + RETURN_IF_ERROR(convert_to_arrow_type(element_doris_type, &field_type, timezone, + datetime_naive)); + fields.push_back(create_arrow_field_with_metadata( + type_struct->get_element_name(i), field_type, element_doris_type->is_nullable(), + element_doris_type->get_primitive_type())); } *result = std::make_shared(fields); break; @@ -184,22 +198,40 @@ Status convert_to_arrow_type(const DataTypePtr& origin_type, return Status::OK(); } -// Helper function to create an Arrow Field with type metadata if applicable, such as IP types +// The Doris types that Arrow has no equivalent for. Each of them travels as some other Arrow type +// and is then indistinguishable from a column that is natively of that type -- a LARGEINT and a +// STRING both arrive as utf8, an IPV4 and an INT both arrive as int32 -- so the original type is +// the only thing that lets a client tell them apart, and it is recorded in the field metadata. +// The names match what the FE reports for the same column under ARROW:FLIGHT:SQL:TYPE_NAME. +static const char* doris_type_metadata_value(PrimitiveType primitive_type) { + switch (primitive_type) { + case PrimitiveType::TYPE_IPV4: + return "IPV4"; + case PrimitiveType::TYPE_IPV6: + return "IPV6"; + case PrimitiveType::TYPE_LARGEINT: + return "LARGEINT"; + case PrimitiveType::TYPE_JSONB: + return "JSON"; + case PrimitiveType::TYPE_VARIANT: + return "VARIANT"; + default: + return nullptr; + } +} + +// Helper function to create an Arrow Field with type metadata if applicable, such as IP types. +// Used for every field of the schema, nested ones included: an element of an ARRAY, MAP or STRUCT +// loses its Doris type in exactly the same way a top level column does. std::shared_ptr create_arrow_field_with_metadata( const std::string& field_name, const std::shared_ptr& arrow_type, bool is_nullable, PrimitiveType primitive_type) { - if (primitive_type == PrimitiveType::TYPE_IPV4) { - auto metadata = arrow::KeyValueMetadata::Make({"doris_type"}, {"IPV4"}); - return std::make_shared(field_name, arrow_type, is_nullable, metadata); - } else if (primitive_type == PrimitiveType::TYPE_IPV6) { - auto metadata = arrow::KeyValueMetadata::Make({"doris_type"}, {"IPV6"}); - return std::make_shared(field_name, arrow_type, is_nullable, metadata); - } else if (primitive_type == PrimitiveType::TYPE_LARGEINT) { - auto metadata = arrow::KeyValueMetadata::Make({"doris_type"}, {"LARGEINT"}); - return std::make_shared(field_name, arrow_type, is_nullable, metadata); - } else { + const char* doris_type = doris_type_metadata_value(primitive_type); + if (doris_type == nullptr) { return std::make_shared(field_name, arrow_type, is_nullable); } + auto metadata = arrow::KeyValueMetadata::Make({"doris_type"}, {doris_type}); + return std::make_shared(field_name, arrow_type, is_nullable, metadata); } Status get_arrow_schema_from_block(const Block& block, std::shared_ptr* result, diff --git a/be/test/format/arrow/arrow_row_batch_test.cpp b/be/test/format/arrow/arrow_row_batch_test.cpp new file mode 100644 index 00000000000000..c3898a328e9473 --- /dev/null +++ b/be/test/format/arrow/arrow_row_batch_test.cpp @@ -0,0 +1,284 @@ +// 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. + +// ############################################################################ +// What an Arrow schema says about the Doris types that Arrow cannot express. +// +// LARGEINT, IPV4, IPV6, JSON and VARIANT all travel as some other Arrow type, +// and once they arrive they are indistinguishable from a column that is +// natively of that type: a LARGEINT and a STRING are both utf8, an IPV4 and an +// INT are both int32. The field metadata is the only thing that tells them +// apart, so a field that loses it loses the type -- and an element of an +// ARRAY, MAP or STRUCT is a field like any other. +// ############################################################################ + +#include "format/arrow/arrow_row_batch.h" + +#include +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include + +#include "core/block/block.h" +#include "core/column/column_array.h" +#include "core/column/column_map.h" +#include "core/column/column_struct.h" +#include "core/data_type/data_type_array.h" +#include "core/data_type/data_type_ipv4.h" +#include "core/data_type/data_type_ipv6.h" +#include "core/data_type/data_type_jsonb.h" +#include "core/data_type/data_type_map.h" +#include "core/data_type/data_type_nullable.h" +#include "core/data_type/data_type_number.h" +#include "core/data_type/data_type_string.h" +#include "core/data_type/data_type_struct.h" +#include "core/data_type/data_type_variant.h" +#include "core/field.h" +#include "format/arrow/arrow_block_convertor.h" + +namespace doris { + +namespace { + +// "" when the field carries no Doris type, which is the answer for every type Arrow can express. +std::string doris_type_of(const std::shared_ptr& field) { + if (field == nullptr || !field->HasMetadata()) { + return ""; + } + const auto found = field->metadata()->Get("doris_type"); + return found.ok() ? found.ValueUnsafe() : ""; +} + +DataTypePtr largeint() { + return make_nullable(std::make_shared()); +} + +DataTypePtr string_type() { + return make_nullable(std::make_shared()); +} + +DataTypePtr array_of(const DataTypePtr& item) { + return make_nullable(std::make_shared(item)); +} + +DataTypePtr map_of(const DataTypePtr& key, const DataTypePtr& value) { + return make_nullable(std::make_shared(key, value)); +} + +DataTypePtr struct_of(const DataTypes& elements, const Strings& names) { + return make_nullable(std::make_shared(elements, names)); +} + +Block block_of(const std::vector>& columns) { + Block block; + for (const auto& [name, type] : columns) { + block.insert(ColumnWithTypeAndName(type->create_column(), type, name)); + } + return block; +} + +std::shared_ptr schema_of( + const std::vector>& columns) { + const Block block = block_of(columns); + std::shared_ptr schema; + EXPECT_TRUE(get_arrow_schema_from_block(block, &schema, "UTC").ok()); + EXPECT_NE(schema, nullptr); + return schema; +} + +// The single child of a list, or the named child of a struct. +std::shared_ptr item_of(const std::shared_ptr& list_field) { + return list_field->type()->field(0); +} + +const arrow::MapType& map_of(const std::shared_ptr& map_field) { + return dynamic_cast(*map_field->type()); +} + +} // namespace + +TEST(ArrowRowBatchSchemaTest, NestedLargeintKeepsItsDorisType) { + auto schema = schema_of({ + {"scalar_value", largeint()}, + {"array_value", array_of(largeint())}, + {"struct_value", struct_of({largeint()}, {"count"})}, + {"map_value", map_of(string_type(), largeint())}, + }); + + // The top level field is the one that already worked. + EXPECT_EQ("LARGEINT", doris_type_of(schema->GetFieldByName("scalar_value"))); + + EXPECT_EQ("LARGEINT", doris_type_of(item_of(schema->GetFieldByName("array_value")))); + EXPECT_EQ("LARGEINT", doris_type_of(schema->GetFieldByName("struct_value")->type()->field(0))); + EXPECT_EQ("LARGEINT", doris_type_of(map_of(schema->GetFieldByName("map_value")).item_field())); +} + +TEST(ArrowRowBatchSchemaTest, LargeintMapKeyKeepsItsDorisType) { + auto schema = schema_of({{"map_value", map_of(largeint(), string_type())}}); + + const auto& map = map_of(schema->GetFieldByName("map_value")); + EXPECT_EQ("LARGEINT", doris_type_of(map.key_field())); + EXPECT_EQ("", doris_type_of(map.item_field())); +} + +TEST(ArrowRowBatchSchemaTest, DorisTypeSurvivesEveryLevelOfNesting) { + auto schema = schema_of({ + {"a", array_of(struct_of({largeint()}, {"n"}))}, + {"b", map_of(string_type(), array_of(array_of(largeint())))}, + }); + + // array> + EXPECT_EQ("LARGEINT", doris_type_of(item_of(schema->GetFieldByName("a"))->type()->field(0))); + + // map>> + const auto& map = map_of(schema->GetFieldByName("b")); + EXPECT_EQ("LARGEINT", doris_type_of(map.item_field()->type()->field(0)->type()->field(0))); +} + +TEST(ArrowRowBatchSchemaTest, NestedIpJsonAndVariantKeepTheirDorisType) { + auto ipv4 = make_nullable(std::make_shared()); + auto ipv6 = make_nullable(std::make_shared()); + auto json = make_nullable(std::make_shared()); + auto variant = make_nullable(std::make_shared()); + + auto schema = schema_of({ + {"ip4", ipv4}, + {"json", json}, + {"variant", variant}, + {"nested", struct_of({ipv4, ipv6, json, variant}, {"ip4", "ip6", "json", "variant"})}, + }); + + // IPV4 is the one whose value changes form -- it arrives as the address's 32 bits read as a + // signed int32 -- so a nested IPV4 that lost its metadata could not be recovered at all. + EXPECT_EQ("IPV4", doris_type_of(schema->GetFieldByName("ip4"))); + EXPECT_EQ("JSON", doris_type_of(schema->GetFieldByName("json"))); + EXPECT_EQ("VARIANT", doris_type_of(schema->GetFieldByName("variant"))); + + const auto& nested = *schema->GetFieldByName("nested")->type(); + EXPECT_EQ("IPV4", doris_type_of(nested.field(0))); + EXPECT_EQ("IPV6", doris_type_of(nested.field(1))); + EXPECT_EQ("JSON", doris_type_of(nested.field(2))); + EXPECT_EQ("VARIANT", doris_type_of(nested.field(3))); +} + +TEST(ArrowRowBatchSchemaTest, TypesArrowCanExpressCarryNoDorisType) { + // The negative half of the contract: if every utf8 field claimed a Doris type there would be + // nothing to distinguish a LARGEINT from a column that really is a string. + auto schema = schema_of({ + {"s", string_type()}, + {"i", make_nullable(std::make_shared())}, + {"nested", struct_of({string_type(), make_nullable(std::make_shared())}, + {"s", "i"})}, + {"arr", array_of(string_type())}, + }); + + EXPECT_EQ("", doris_type_of(schema->GetFieldByName("s"))); + EXPECT_EQ("", doris_type_of(schema->GetFieldByName("i"))); + EXPECT_EQ("", doris_type_of(schema->GetFieldByName("nested")->type()->field(0))); + EXPECT_EQ("", doris_type_of(schema->GetFieldByName("nested")->type()->field(1))); + EXPECT_EQ("", doris_type_of(item_of(schema->GetFieldByName("arr")))); +} + +TEST(ArrowRowBatchSchemaTest, OnlyTheMetadataIsNew) { + // Naming a child or changing its nullability would be a different schema, and clients that + // already read these columns type themselves from it. Compared against the types Arrow's own + // constructors build: equal when metadata is ignored, different only once it is compared. + auto schema = schema_of({ + {"array_value", array_of(largeint())}, + {"map_value", map_of(string_type(), largeint())}, + }); + + const auto expected_list = arrow::list(arrow::utf8()); + const auto& list = *schema->GetFieldByName("array_value")->type(); + EXPECT_TRUE(list.Equals(*expected_list, /*check_metadata=*/false)); + EXPECT_FALSE(list.Equals(*expected_list, /*check_metadata=*/true)); + EXPECT_EQ("item", list.field(0)->name()); + EXPECT_TRUE(list.field(0)->nullable()); + + const auto expected_map = arrow::map(arrow::utf8(), arrow::utf8()); + const auto& map = *schema->GetFieldByName("map_value")->type(); + EXPECT_TRUE(map.Equals(*expected_map, /*check_metadata=*/false)); + EXPECT_FALSE(map.Equals(*expected_map, /*check_metadata=*/true)); + const auto& as_map = dynamic_cast(map); + EXPECT_EQ("key", as_map.key_field()->name()); + EXPECT_FALSE(as_map.key_field()->nullable()) << "an Arrow map key is never nullable"; + EXPECT_EQ("value", as_map.item_field()->name()); + EXPECT_TRUE(as_map.item_field()->nullable()); +} + +// The schema is not only described to the client, it is also what the record batch builders are +// made from (FromBlockToRecordBatchConverter reads _schema->field(idx)->type()). A schema the data +// path cannot honour would turn a metadata fix into a broken result set. +TEST(ArrowRowBatchDataTest, NestedSchemaStillBuildsTheBatch) { + const __int128_t value = 495; + + auto array_type = array_of(largeint()); + auto struct_type = struct_of({largeint()}, {"count"}); + auto map_type = map_of(string_type(), largeint()); + + auto array_column = array_type->create_column(); + array_column->insert( + Field::create_field(Array {Field::create_field(value)})); + + auto struct_column = struct_type->create_column(); + struct_column->insert( + Field::create_field(Struct {Field::create_field(value)})); + + auto map_column = map_type->create_column(); + map_column->insert(Field::create_field(Map { + Field::create_field(Array {Field::create_field(String("k"))}), + Field::create_field(Array {Field::create_field(value)})})); + + Block block; + block.insert(ColumnWithTypeAndName(std::move(array_column), array_type, "array_value")); + block.insert(ColumnWithTypeAndName(std::move(struct_column), struct_type, "struct_value")); + block.insert(ColumnWithTypeAndName(std::move(map_column), map_type, "map_value")); + + std::shared_ptr schema; + ASSERT_TRUE(get_arrow_schema_from_block(block, &schema, "UTC").ok()); + + std::shared_ptr batch; + cctz::time_zone utc; + ASSERT_TRUE( + convert_to_arrow_batch(block, schema, arrow::default_memory_pool(), &batch, utc).ok()); + ASSERT_NE(batch, nullptr); + ASSERT_TRUE(batch->ValidateFull().ok()) << batch->ValidateFull().ToString(); + // Including the metadata: the batch the client reads describes its nested types the same way + // the schema does. + EXPECT_TRUE(batch->schema()->Equals(*schema, /*check_metadata=*/true)); + ASSERT_EQ(1, batch->num_rows()); + + const auto& list = dynamic_cast(*batch->column(0)); + EXPECT_EQ("495", dynamic_cast(*list.values()).GetString(0)); + + const auto& structure = dynamic_cast(*batch->column(1)); + EXPECT_EQ("495", dynamic_cast(*structure.field(0)).GetString(0)); + + const auto& map = dynamic_cast(*batch->column(2)); + EXPECT_EQ("k", dynamic_cast(*map.keys()).GetString(0)); + EXPECT_EQ("495", dynamic_cast(*map.items()).GetString(0)); +} + +} // namespace doris