diff --git a/docs/source/user_guide/compaction.rst b/docs/source/user_guide/compaction.rst index 492e27ff..899a74f1 100644 --- a/docs/source/user_guide/compaction.rst +++ b/docs/source/user_guide/compaction.rst @@ -118,9 +118,10 @@ one is read as ordinary values, and the writer takes both. A high-cardinality column that started dictionary-encoded and fell back to plain encoding therefore does not qualify, even though it still carries a dictionary page. ``BINARY`` is not forwarded although Parquet stores it in the same physical type and -dictionary-encodes it the same way, because the value accessors cannot read a -``BINARY`` dictionary. The rewrite also overrides the option back to ``false`` -when the table writes a format other than Parquet, when +dictionary-encodes it the same way; this option only selects ``STRING`` columns. +A file with a serialized Arrow dictionary schema can independently produce a +``BINARY`` dictionary, which the value accessors also support. The rewrite also +overrides the option back to ``false`` when the table writes a format other than Parquet, when ``parquet.enable-dictionary`` is ``false`` because the writer would only expand the values again, or when variant/map shredding is configured because those writers reshape each batch against a fixed physical schema. Setting the option diff --git a/src/paimon/common/data/columnar/columnar_utils.h b/src/paimon/common/data/columnar/columnar_utils.h index ede3155b..2ab413e3 100644 --- a/src/paimon/common/data/columnar/columnar_utils.h +++ b/src/paimon/common/data/columnar/columnar_utils.h @@ -73,13 +73,15 @@ class ColumnarUtils { dict_index = indices->Value(pos); } assert(dict_index >= 0); - if (value_type_id == arrow::Type::type::STRING) { + if (value_type_id == arrow::Type::type::STRING || + value_type_id == arrow::Type::type::BINARY) { auto dictionary = - checked_cast(typed_array->dictionary().get()); + checked_cast(typed_array->dictionary().get()); return dictionary->GetView(dict_index); - } else if (value_type_id == arrow::Type::type::LARGE_STRING) { + } else if (value_type_id == arrow::Type::type::LARGE_STRING || + value_type_id == arrow::Type::type::LARGE_BINARY) { auto dictionary = - checked_cast(typed_array->dictionary().get()); + checked_cast(typed_array->dictionary().get()); return dictionary->GetView(dict_index); } assert(false); diff --git a/src/paimon/common/data/columnar/columnar_utils_test.cpp b/src/paimon/common/data/columnar/columnar_utils_test.cpp index 7eb81879..4b6449a5 100644 --- a/src/paimon/common/data/columnar/columnar_utils_test.cpp +++ b/src/paimon/common/data/columnar/columnar_utils_test.cpp @@ -18,6 +18,7 @@ #include "paimon/common/data/columnar/columnar_utils.h" #include +#include #include "arrow/api.h" #include "arrow/array/array_dict.h" @@ -53,4 +54,51 @@ TEST(ColumnarUtilsTest, TestGetViewAndBytesOfDict) { ASSERT_EQ("foo", std::string(ColumnarUtils::GetView(dict_array.get(), 4))); } +template +class ColumnarUtilsBinaryDictionaryTest : public ::testing::Test {}; + +using BinaryDictionaryTypes = ::testing::Types; +TYPED_TEST_SUITE(ColumnarUtilsBinaryDictionaryTest, BinaryDictionaryTypes); + +TYPED_TEST(ColumnarUtilsBinaryDictionaryTest, GetViewAndBytes) { + auto pool = GetDefaultPool(); + const std::vector values = {std::string("\x00\xff\x80", 3), "", + std::string("a\0b", 3)}; + typename arrow::TypeTraits::BuilderType builder; + ASSERT_TRUE(builder.Append("unused").ok()); + for (const auto& value : values) { + ASSERT_TRUE(builder.Append(value).ok()); + } + std::shared_ptr dictionary; + ASSERT_TRUE(builder.Finish(&dictionary).ok()); + dictionary = dictionary->Slice(1); + + const std::vector> index_types = { + arrow::int8(), arrow::int16(), arrow::int32(), arrow::int64()}; + const std::vector expected_indices = {0, 1, 2, 0, -1, 2}; + for (const auto& index_type : index_types) { + SCOPED_TRACE(index_type->ToString()); + auto indices = + arrow::ipc::internal::json::ArrayFromJSON(index_type, "[0, 1, 2, 0, null, 2]") + .ValueOrDie(); + auto dict_array = arrow::DictionaryArray::FromArrays(indices, dictionary).ValueOrDie(); + ASSERT_TRUE(dict_array->ValidateFull().ok()); + for (int32_t offset : {0, 1}) { + SCOPED_TRACE(offset); + auto sliced = dict_array->Slice(offset); + for (int32_t pos = 0; pos < sliced->length(); ++pos) { + SCOPED_TRACE(pos); + int32_t index = expected_indices[offset + pos]; + ASSERT_EQ(index == -1, sliced->IsNull(pos)); + if (index == -1) { + continue; + } + ASSERT_EQ(values[index], ColumnarUtils::GetView(sliced.get(), pos)); + auto bytes = ColumnarUtils::GetBytes(sliced.get(), pos, pool.get()); + ASSERT_EQ(Bytes(values[index], pool.get()), *bytes); + } + } + } +} + } // namespace paimon::test diff --git a/src/paimon/common/predicate/literal_converter.cpp b/src/paimon/common/predicate/literal_converter.cpp index 4129fe70..2b1a1eb2 100644 --- a/src/paimon/common/predicate/literal_converter.cpp +++ b/src/paimon/common/predicate/literal_converter.cpp @@ -214,17 +214,26 @@ Result> LiteralConverter::ConvertLiteralsFromArray(const ar auto* dict_type = checked_cast(dict_array.type().get()); auto value_type_id = dict_type->value_type()->id(); auto index_type_id = dict_type->index_type()->id(); - if (value_type_id == arrow::Type::type::STRING && + if ((value_type_id == arrow::Type::type::STRING || + value_type_id == arrow::Type::type::BINARY) && index_type_id == arrow::Type::type::INT32) { - return GetLiteralFromDictionaryArray( - dict_array, FieldType::STRING, own_data); - } else if (value_type_id == arrow::Type::type::LARGE_STRING && + FieldType literal_type = value_type_id == arrow::Type::type::STRING + ? FieldType::STRING + : FieldType::BINARY; + return GetLiteralFromDictionaryArray( + dict_array, literal_type, own_data); + } else if ((value_type_id == arrow::Type::type::LARGE_STRING || + value_type_id == arrow::Type::type::LARGE_BINARY) && index_type_id == arrow::Type::type::INT64) { - return GetLiteralFromDictionaryArray( - dict_array, FieldType::STRING, own_data); + FieldType literal_type = value_type_id == arrow::Type::type::LARGE_STRING + ? FieldType::STRING + : FieldType::BINARY; + return GetLiteralFromDictionaryArray( + dict_array, literal_type, own_data); } else { return Status::Invalid( - "only support [STRING, INT32] or [LARGE_STRING, INT64] for DictionaryArray"); + "only support [STRING|BINARY, INT32] or " + "[LARGE_STRING|LARGE_BINARY, INT64] for DictionaryArray"); } } default: diff --git a/src/paimon/common/predicate/literal_converter.h b/src/paimon/common/predicate/literal_converter.h index a3fde588..7133812f 100644 --- a/src/paimon/common/predicate/literal_converter.h +++ b/src/paimon/common/predicate/literal_converter.h @@ -134,17 +134,12 @@ class PAIMON_EXPORT LiteralConverter { literals.emplace_back(literal_type); } else { int64_t dict_index = indices->Value(i); - if constexpr (std::is_same_v) { - int32_t length = 0; - const uint8_t* value = dictionary->GetValue(dict_index, &length); - literals.emplace_back(literal_type, reinterpret_cast(value), - length, own_data); - } else { - int64_t length = 0; - const uint8_t* value = dictionary->GetValue(dict_index, &length); - literals.emplace_back(literal_type, reinterpret_cast(value), - length, own_data); + if (dictionary->IsNull(dict_index)) { + literals.emplace_back(literal_type); + continue; } + auto value = dictionary->GetView(dict_index); + literals.emplace_back(literal_type, value.data(), value.size(), own_data); } } return literals; diff --git a/src/paimon/common/predicate/literal_converter_test.cpp b/src/paimon/common/predicate/literal_converter_test.cpp index 0d36cfe7..ccdce40f 100644 --- a/src/paimon/common/predicate/literal_converter_test.cpp +++ b/src/paimon/common/predicate/literal_converter_test.cpp @@ -449,6 +449,59 @@ TEST_F(LiteralConverterTest, TestDictType) { Literal(FieldType::STRING, "foo", 3), Literal(FieldType::STRING)})); } +TEST_F(LiteralConverterTest, TestDictionaryValueAndIndexTypes) { + const std::vector, std::shared_ptr>> + types = {{arrow::utf8(), arrow::int32()}, + {arrow::binary(), arrow::int32()}, + {arrow::large_utf8(), arrow::int64()}, + {arrow::large_binary(), arrow::int64()}}; + for (const auto& [value_type, index_type] : types) { + SCOPED_TRACE(value_type->ToString()); + const FieldType literal_type = + arrow::is_string(value_type->id()) ? FieldType::STRING : FieldType::BINARY; + auto dictionary = arrow::ipc::internal::json::ArrayFromJSON( + value_type, R"(["unused", "a\u0000b", "", null, "tail"])") + .ValueOrDie() + ->Slice(1); + const std::vector expected = { + Literal(literal_type, "a\0b", 3), Literal(literal_type), + Literal(literal_type, "", 0), Literal(literal_type), + Literal(literal_type, "tail", 4), Literal(literal_type, "a\0b", 3)}; + SCOPED_TRACE(index_type->ToString()); + auto indices = + arrow::ipc::internal::json::ArrayFromJSON(index_type, "[3, 0, null, 1, 2, 3, 0]") + .ValueOrDie(); + auto array = arrow::DictionaryArray::FromArrays(indices, dictionary).ValueOrDie(); + ASSERT_TRUE(array->ValidateFull().ok()); + auto sliced = array->Slice(1); + for (bool own_data : {false, true}) { + SCOPED_TRACE(own_data); + ASSERT_OK_AND_ASSIGN(std::vector actual, + LiteralConverter::ConvertLiteralsFromArray(*sliced, own_data)); + ASSERT_EQ(expected, actual); + } + } +} + +TEST_F(LiteralConverterTest, TestBinaryDictionaryPreservesBytes) { + const std::string bytes("\x00\xff\x80", 3); + arrow::BinaryBuilder builder; + ASSERT_TRUE(builder.Append(bytes).ok()); + std::shared_ptr dictionary; + ASSERT_TRUE(builder.Finish(&dictionary).ok()); + auto indices = + arrow::ipc::internal::json::ArrayFromJSON(arrow::int32(), "[0, null, 0]").ValueOrDie(); + auto array = arrow::DictionaryArray::FromArrays(indices, dictionary).ValueOrDie(); + ASSERT_OK_AND_ASSIGN(std::vector actual, + LiteralConverter::ConvertLiteralsFromArray(*array, /*own_data=*/true)); + array.reset(); + dictionary.reset(); + ASSERT_EQ(actual, + std::vector({Literal(FieldType::BINARY, bytes.data(), bytes.size()), + Literal(FieldType::BINARY), + Literal(FieldType::BINARY, bytes.data(), bytes.size())})); +} + TEST_F(LiteralConverterTest, TestLiteralsToArray) { // Every writable field type converts to an array and back, the null literal included, so each // case asserts `ConvertLiteralsToArray` and `ConvertLiteralsFromArray` agree on every type. diff --git a/src/paimon/common/predicate/null_false_leaf_binary_function.cpp b/src/paimon/common/predicate/null_false_leaf_binary_function.cpp index 5968ca28..7aa95759 100644 --- a/src/paimon/common/predicate/null_false_leaf_binary_function.cpp +++ b/src/paimon/common/predicate/null_false_leaf_binary_function.cpp @@ -98,12 +98,15 @@ std::optional ArrayFieldType(const arrow::Array& array) { const std::shared_ptr& type = array.type(); if (type->id() == arrow::Type::DICTIONARY) { const auto& dict_type = checked_cast(*type); - const bool is_string = dict_type.value_type()->id() == arrow::Type::STRING && - dict_type.index_type()->id() == arrow::Type::INT32; - const bool is_large_string = dict_type.value_type()->id() == arrow::Type::LARGE_STRING && - dict_type.index_type()->id() == arrow::Type::INT64; - if (is_string || is_large_string) { - return FieldType::STRING; + const arrow::Type::type value_type = dict_type.value_type()->id(); + const arrow::Type::type index_type = dict_type.index_type()->id(); + if ((value_type == arrow::Type::STRING || value_type == arrow::Type::BINARY) && + index_type == arrow::Type::INT32) { + return value_type == arrow::Type::STRING ? FieldType::STRING : FieldType::BINARY; + } + if ((value_type == arrow::Type::LARGE_STRING || value_type == arrow::Type::LARGE_BINARY) && + index_type == arrow::Type::INT64) { + return value_type == arrow::Type::LARGE_STRING ? FieldType::STRING : FieldType::BINARY; } return std::nullopt; } diff --git a/src/paimon/common/predicate/null_false_leaf_binary_function_test.cpp b/src/paimon/common/predicate/null_false_leaf_binary_function_test.cpp index ff23eafd..b757a460 100644 --- a/src/paimon/common/predicate/null_false_leaf_binary_function_test.cpp +++ b/src/paimon/common/predicate/null_false_leaf_binary_function_test.cpp @@ -339,6 +339,28 @@ TEST_F(NullFalseLeafBinaryFunctionTest, TestLargeStringDictionary) { std::vector({0, 1, 0, 0})); } +TEST_F(NullFalseLeafBinaryFunctionTest, TestBinaryDictionary) { + const std::vector, std::shared_ptr>> + types = {{arrow::binary(), arrow::int32()}, {arrow::large_binary(), arrow::int64()}}; + for (const auto& [value_type, index_type] : types) { + SCOPED_TRACE(value_type->ToString()); + auto dictionary = + arrow::ipc::internal::json::ArrayFromJSON(value_type, R"(["a\u0000b", "a", "", null])") + .ValueOrDie(); + SCOPED_TRACE(index_type->ToString()); + auto indices = + arrow::ipc::internal::json::ArrayFromJSON(index_type, "[1, 0, null, 1, 2, 3, 0]") + .ValueOrDie(); + auto array = arrow::DictionaryArray::FromArrays(indices, dictionary).ValueOrDie()->Slice(1); + const auto literal = BinaryLiteral(std::string("a\0b", 3)); + ASSERT_EQ(Eval(Equal::Instance(), literal, array), std::vector({1, 0, 0, 0, 0, 1})); + ASSERT_EQ(Eval(LessThan::Instance(), literal, array), + std::vector({0, 0, 1, 1, 0, 0})); + ASSERT_EQ(Eval(NotEqual::Instance(), BinaryLiteral(""), array), + std::vector({1, 0, 1, 0, 0, 1})); + } +} + TEST_F(NullFalseLeafBinaryFunctionTest, TestDictionaryWithNullValue) { // A kernel decodes the dictionary, so a row pointing at a null dictionary value becomes a null // row and is false for every comparison. The row by row path reads the slot the index names diff --git a/src/paimon/core/casting/casting_utils.cpp b/src/paimon/core/casting/casting_utils.cpp index d17b475c..90e6998a 100644 --- a/src/paimon/core/casting/casting_utils.cpp +++ b/src/paimon/core/casting/casting_utils.cpp @@ -23,6 +23,19 @@ #include "paimon/common/utils/checked_cast.h" namespace paimon { +Result> CastingUtils::DecodeDictionary( + const std::shared_ptr& array, arrow::MemoryPool* pool) { + if (array->type_id() != arrow::Type::DICTIONARY) { + return array; + } + const auto& dictionary_type = checked_cast(*array->type()); + std::shared_ptr value_type = dictionary_type.value_type(); + if (value_type->id() == arrow::Type::LARGE_STRING) { + value_type = arrow::utf8(); + } + return Cast(array, value_type, arrow::compute::CastOptions::Safe(), pool); +} + Result> CastingUtils::Cast( const std::shared_ptr& src_array, const std::shared_ptr& target_type, const arrow::compute::CastOptions& options, diff --git a/src/paimon/core/casting/casting_utils.h b/src/paimon/core/casting/casting_utils.h index 8d27f06c..2c34107f 100644 --- a/src/paimon/core/casting/casting_utils.h +++ b/src/paimon/core/casting/casting_utils.h @@ -47,6 +47,12 @@ class PAIMON_EXPORT CastingUtils { const std::shared_ptr& target_type, const arrow::compute::CastOptions& options, arrow::MemoryPool* pool); + /// Decodes a dictionary to its value type, normalizing LARGE_STRING to STRING for readers + /// that widen string offsets. Other value types, including binary offsets, are preserved. + /// Returns a non-dictionary array unchanged. + static Result> DecodeDictionary( + const std::shared_ptr& array, arrow::MemoryPool* pool); + template static Result Cast(const Literal& literal, diff --git a/src/paimon/core/casting/casting_utils_test.cpp b/src/paimon/core/casting/casting_utils_test.cpp index 0c66c662..ec78d261 100644 --- a/src/paimon/core/casting/casting_utils_test.cpp +++ b/src/paimon/core/casting/casting_utils_test.cpp @@ -53,6 +53,59 @@ TEST_F(CastingUtilsTest, TestDictionaryToString) { ASSERT_TRUE(result_array->Equals(string_array)); } +TEST_F(CastingUtilsTest, TestDecodeDictionaryPreservesValueType) { + for (const auto& value_type : {arrow::utf8(), arrow::large_utf8(), arrow::binary(), + arrow::large_binary(), arrow::int32()}) { + SCOPED_TRACE(value_type->ToString()); + const bool is_integer = value_type->id() == arrow::Type::INT32; + auto dictionary = + arrow::ipc::internal::json::ArrayFromJSON( + value_type, is_integer ? "[10, null, 20]" : R"(["a\u0000b", null, ""])") + .ValueOrDie(); + const auto target_type = + value_type->id() == arrow::Type::LARGE_STRING ? arrow::utf8() : value_type; + auto expected = arrow::ipc::internal::json::ArrayFromJSON( + target_type, is_integer ? "[10, null, null, 20, 10]" + : R"(["a\u0000b", null, null, "", "a\u0000b"])") + .ValueOrDie(); + for (const auto& index_type : + {arrow::int8(), arrow::int16(), arrow::int32(), arrow::int64()}) { + SCOPED_TRACE(index_type->ToString()); + auto indices = + arrow::ipc::internal::json::ArrayFromJSON(index_type, "[2, 0, null, 1, 2, 0]") + .ValueOrDie(); + auto array = arrow::DictionaryArray::FromArrays(indices, dictionary).ValueOrDie(); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr decoded, + CastingUtils::DecodeDictionary(array->Slice(1), arrow_pool_.get())); + ASSERT_TRUE(decoded->ValidateFull().ok()); + ASSERT_TRUE(decoded->Equals(expected)) << decoded->ToString(); + } + ASSERT_OK_AND_ASSIGN(std::shared_ptr unchanged, + CastingUtils::DecodeDictionary(dictionary, arrow_pool_.get())); + ASSERT_EQ(unchanged, dictionary); + } +} + +TEST_F(CastingUtilsTest, TestDecodeBinaryDictionaryPreservesNonUtf8) { + const std::string bytes("\x00\xff\x80", 3); + arrow::BinaryBuilder builder; + ASSERT_TRUE(builder.Append(bytes).ok()); + std::shared_ptr dictionary; + ASSERT_TRUE(builder.Finish(&dictionary).ok()); + auto indices = + arrow::ipc::internal::json::ArrayFromJSON(arrow::int32(), "[0, null, 0]").ValueOrDie(); + auto array = arrow::DictionaryArray::FromArrays(indices, dictionary).ValueOrDie(); + ASSERT_OK_AND_ASSIGN(std::shared_ptr decoded, + CastingUtils::DecodeDictionary(array, arrow_pool_.get())); + ASSERT_TRUE(decoded->ValidateFull().ok()); + ASSERT_EQ(decoded->type_id(), arrow::Type::BINARY); + auto binary = checked_pointer_cast(decoded); + ASSERT_EQ(bytes, binary->GetView(0)); + ASSERT_TRUE(binary->IsNull(1)); + ASSERT_EQ(bytes, binary->GetView(2)); +} + TEST_F(CastingUtilsTest, TestTimestampToTimestampWithTimezone) { // local no tz -> utc tz auto src_array = arrow::ipc::internal::json::ArrayFromJSON( diff --git a/src/paimon/core/global_index/global_index_write_task.cpp b/src/paimon/core/global_index/global_index_write_task.cpp index 0288d1b9..89aea4a1 100644 --- a/src/paimon/core/global_index/global_index_write_task.cpp +++ b/src/paimon/core/global_index/global_index_write_task.cpp @@ -170,22 +170,11 @@ Result> CreateBatchReader( return table_read->CreateReader(indexed_split); } -Result> CastDictionaryArrayToString( +Result> DecodeDictionaryArrays( const std::shared_ptr& array, arrow::MemoryPool* pool) { arrow::Type::type type_id = array->type_id(); if (type_id == arrow::Type::DICTIONARY) { - const auto* dictionary_type = - checked_cast(array->type().get()); - arrow::Type::type value_type = dictionary_type->value_type()->id(); - if (value_type != arrow::Type::STRING && value_type != arrow::Type::LARGE_STRING) { - return Status::Invalid(fmt::format( - "GlobalIndexWriteTask cannot decode dictionary array with value type {}", - dictionary_type->value_type()->ToString())); - } - PAIMON_ASSIGN_OR_RAISE( - std::shared_ptr casted_array, - CastingUtils::Cast(array, arrow::utf8(), arrow::compute::CastOptions::Safe(), pool)); - return casted_array; + return CastingUtils::DecodeDictionary(array, pool); } if (type_id != arrow::Type::STRUCT && type_id != arrow::Type::MAP && type_id != arrow::Type::LIST) { @@ -199,7 +188,7 @@ Result> CastDictionaryArrayToString( for (int32_t i = 0; i < struct_array->num_fields(); i++) { std::shared_ptr child = struct_array->field(i); PAIMON_ASSIGN_OR_RAISE(std::shared_ptr casted_child, - CastDictionaryArrayToString(child, pool)); + DecodeDictionaryArrays(child, pool)); if (casted_child != child && children.empty()) { children = struct_array->fields(); } @@ -227,9 +216,9 @@ Result> CastDictionaryArrayToString( std::shared_ptr original_keys = map_array->keys(); std::shared_ptr original_items = map_array->items(); PAIMON_ASSIGN_OR_RAISE(std::shared_ptr keys, - CastDictionaryArrayToString(original_keys, pool)); + DecodeDictionaryArrays(original_keys, pool)); PAIMON_ASSIGN_OR_RAISE(std::shared_ptr items, - CastDictionaryArrayToString(original_items, pool)); + DecodeDictionaryArrays(original_items, pool)); if (keys == original_keys && items == original_items) { return array; } @@ -245,7 +234,7 @@ Result> CastDictionaryArrayToString( std::shared_ptr list_array = checked_pointer_cast(array); std::shared_ptr original_values = list_array->values(); PAIMON_ASSIGN_OR_RAISE(std::shared_ptr values, - CastDictionaryArrayToString(original_values, pool)); + DecodeDictionaryArrays(original_values, pool)); if (values == original_values) { return array; } @@ -302,7 +291,7 @@ Result> BuildIndex( writer_field_name)); } PAIMON_ASSIGN_OR_RAISE(std::shared_ptr decoded_writer_array, - CastDictionaryArrayToString(writer_array, arrow_pool)); + DecodeDictionaryArrays(writer_array, arrow_pool)); writer_arrays.push_back(std::move(decoded_writer_array)); } PAIMON_ASSIGN_OR_RAISE_FROM_ARROW( diff --git a/src/paimon/core/io/field_mapping_reader.cpp b/src/paimon/core/io/field_mapping_reader.cpp index e2c8b4cf..c7f6ba22 100644 --- a/src/paimon/core/io/field_mapping_reader.cpp +++ b/src/paimon/core/io/field_mapping_reader.cpp @@ -132,24 +132,17 @@ Result> FieldMappingReader::CastNonPartitionArrayI for (int32_t i = 0; i < field_count; i++) { std::shared_ptr column; if (non_partition_info_.cast_executors[i] != nullptr) { - auto single_column_array = struct_array->field(i); - // if src array is dict, cast to string first - auto dict_array = - std::dynamic_pointer_cast(single_column_array); - if (dict_array) { - PAIMON_ASSIGN_OR_RAISE( - single_column_array, - CastingUtils::Cast(dict_array, /*target_type=*/arrow::utf8(), - arrow::compute::CastOptions::Safe(), arrow_pool_.get())); - } + PAIMON_ASSIGN_OR_RAISE( + std::shared_ptr single_column_array, + CastingUtils::DecodeDictionary(struct_array->field(i), arrow_pool_.get())); PAIMON_ASSIGN_OR_RAISE( column, non_partition_info_.cast_executors[i]->Cast( single_column_array, non_partition_info_.non_partition_read_schema[i].Type(), arrow_pool_.get())); } else { - // read and data type may both be string type, but after adapter transform, type may be - // dictionary, need reconstruct struct type + // The reader may return a dictionary for an unchanged logical type, so reconstruct + // the struct type from the actual column. column = struct_array->field(i); } // Null-fill nested fields added by schema evolution. Only when the data and diff --git a/src/paimon/core/io/field_mapping_reader_test.cpp b/src/paimon/core/io/field_mapping_reader_test.cpp index 3756475f..45b5da0f 100644 --- a/src/paimon/core/io/field_mapping_reader_test.cpp +++ b/src/paimon/core/io/field_mapping_reader_test.cpp @@ -611,6 +611,50 @@ TEST_F(FieldMappingReaderTest, TestSchemaEvolutionWithModifyTypeWithDict) { partition, expected_array); } +TEST_F(FieldMappingReaderTest, TestBinaryDictionaryWithSchemaEvolution) { + const std::vector data_fields = { + DataField(0, arrow::field("payload", arrow::binary()))}; + auto read_schema = DataField::ConvertDataFieldsToArrowSchema( + {DataField(0, arrow::field("payload", arrow::utf8()))}); + const std::string bytes("a\0\xff\x80", 4); + arrow::BinaryBuilder builder; + ASSERT_TRUE(builder.Append(bytes).ok()); + ASSERT_TRUE(builder.Append("").ok()); + ASSERT_TRUE(builder.AppendNull().ok()); + std::shared_ptr dictionary; + ASSERT_TRUE(builder.Finish(&dictionary).ok()); + auto indices = + arrow::ipc::internal::json::ArrayFromJSON(arrow::int32(), "[1, 0, null, 1, 2, 0]") + .ValueOrDie(); + auto encoded = arrow::DictionaryArray::FromArrays(indices, dictionary).ValueOrDie()->Slice(1); + auto data = arrow::StructArray::Make({encoded}, {"payload"}).ValueOrDie(); + auto file_type = arrow::struct_({arrow::field("payload", arrow::binary())}); + ASSERT_OK_AND_ASSIGN(std::unique_ptr mapping_builder, + FieldMappingBuilder::Create(read_schema, /*partition_keys=*/{}, + /*predicate=*/nullptr)); + ASSERT_OK_AND_ASSIGN(std::unique_ptr mapping, + mapping_builder->CreateFieldMapping(data_fields)); + auto mock = std::make_unique(data, file_type, /*read_batch_size=*/8); + ASSERT_OK_AND_ASSIGN(std::unique_ptr reader, + FieldMappingReader::Create( + 1, std::move(mock), BinaryRow::EmptyRow(), std::move(mapping), + /*skip_map_selected_keys_filter_field_ids=*/{}, GetArrowPool(pool_))); + ASSERT_OK_AND_ASSIGN(std::shared_ptr result, + ReadResultCollector::CollectResult(std::move(reader))); + arrow::StringBuilder expected_builder; + // BINARY -> STRING intentionally permits invalid UTF-8, just as the non-dictionary path does. + ASSERT_TRUE(expected_builder.Append(bytes).ok()); + ASSERT_TRUE(expected_builder.AppendNull().ok()); + ASSERT_TRUE(expected_builder.Append("").ok()); + ASSERT_TRUE(expected_builder.AppendNull().ok()); + ASSERT_TRUE(expected_builder.Append(bytes).ok()); + std::shared_ptr expected_values; + ASSERT_TRUE(expected_builder.Finish(&expected_values).ok()); + auto expected = arrow::StructArray::Make({expected_values}, {"payload"}).ValueOrDie(); + ASSERT_TRUE(result->Equals(arrow::ChunkedArray(arrow::ArrayVector({expected})))) + << result->ToString(); +} + TEST_F(FieldMappingReaderTest, TestSchemaEvolutionWithModifyTypeWithPredicate) { std::vector data_fields = {DataField(0, arrow::field("f0", arrow::utf8())), DataField(1, arrow::field("f1", arrow::float32())), diff --git a/src/paimon/core/utils/nested_projection_utils.cpp b/src/paimon/core/utils/nested_projection_utils.cpp index 44dd43a0..ff6368a7 100644 --- a/src/paimon/core/utils/nested_projection_utils.cpp +++ b/src/paimon/core/utils/nested_projection_utils.cpp @@ -713,10 +713,9 @@ NestedProjectionUtils::FilterMapArrayBySelectedKeysRecursively( namespace { // Strips physical-only differences from a leaf type: ORC lazy decoding wraps -// strings in a dictionary and may widen them to large_string. binary is not -// dictionary-encoded and large_binary is blob's real type, so neither is -// normalized. Two leaves with equal normalized types hold the same logical -// values. +// strings in a dictionary and may widen them to large_string. Unwrap dictionaries +// for every value type, but preserve large_binary because it is blob's real type. +// Two leaves with equal normalized types hold the same logical values. std::shared_ptr NormalizeLeafRepresentation( const std::shared_ptr& type) { auto t = type; diff --git a/src/paimon/format/parquet/parquet_file_batch_reader.cpp b/src/paimon/format/parquet/parquet_file_batch_reader.cpp index faa18c3e..8d5fdc53 100644 --- a/src/paimon/format/parquet/parquet_file_batch_reader.cpp +++ b/src/paimon/format/parquet/parquet_file_batch_reader.cpp @@ -210,15 +210,9 @@ std::set ParquetFileBatchReader::ResolveFullyDictionaryEncodedColumns( // Arrow only reads BYTE_ARRAY leaves as dictionaries, and only a top-level column can be // forwarded to the writer without rebuilding the nesting around it. // - // `is_string()` narrows that further to STRING, leaving out the other BYTE_ARRAY leaf, - // BINARY. This is the reader's restriction, not the format's: the writer takes - // `dictionary(int32, binary)` and ArrowUtils::IsDictionaryLayoutRecoverableValueType() - // accepts it, but the option applies to every read of the table and the value accessors - // cannot read one. ColumnarUtils::GetView() asserts on a dictionary whose values are - // neither STRING nor LARGE_STRING and returns an empty view in a release build, and - // LiteralConverter rejects it. Every consumer here understands a STRING dictionary because - // the ORC reader has always produced one under lazy decoding; none was ever handed a - // BINARY one. Widening this needs those consumers first, not just the gate. + // Keep passthrough opt-in limited to STRING. This is an output policy, not a Parquet + // restriction: Arrow can also return BINARY dictionaries, including when restoring a + // dictionary type from a file's serialized Arrow schema. if (schema->Column(i)->physical_type() == ::parquet::Type::BYTE_ARRAY && schema->Column(i)->logical_type()->is_string() && schema->GetColumnRoot(i)->is_primitive()) { diff --git a/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp b/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp index f25e0527..d8572f49 100644 --- a/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp +++ b/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp @@ -40,10 +40,14 @@ #include "arrow/io/caching.h" #include "arrow/io/file.h" #include "arrow/io/interfaces.h" +#include "arrow/io/memory.h" #include "arrow/ipc/api.h" #include "arrow/ipc/json_simple.h" #include "fmt/format.h" #include "gtest/gtest.h" +#include "paimon/common/data/columnar/columnar_array.h" +#include "paimon/common/data/columnar/columnar_row.h" +#include "paimon/common/data/columnar/columnar_row_ref.h" #include "paimon/common/io/cache_input_stream.h" #include "paimon/common/metrics/metrics_impl.h" #include "paimon/common/types/data_field.h" @@ -735,6 +739,120 @@ TEST_F(ParquetFileBatchReaderTest, TestNextBatchWithDictionary) { check_result(false); } +TEST_F(ParquetFileBatchReaderTest, TestBinaryDictionaryOutputDependsOnStoredArrowSchema) { + auto dictionary = + arrow::ipc::internal::json::ArrayFromJSON(arrow::binary(), R"(["a\u0000b", ""])") + .ValueOrDie(); + auto indices = arrow::ipc::internal::json::ArrayFromJSON(arrow::int32(), "[0, 1, null, 0, 1]") + .ValueOrDie(); + auto encoded = arrow::DictionaryArray::FromArrays(indices, dictionary).ValueOrDie(); + auto schema = arrow::schema({arrow::field("payload", encoded->type())}); + auto table = arrow::Table::Make(schema, {encoded}); + for (bool store_schema : {false, true}) { + SCOPED_TRACE(store_schema); + auto sink = arrow::io::BufferOutputStream::Create().ValueOrDie(); + auto writer_properties = + ::parquet::WriterProperties::Builder().enable_dictionary()->build(); + ::parquet::ArrowWriterProperties::Builder arrow_properties; + if (store_schema) { + arrow_properties.store_schema(); + } + ASSERT_TRUE(::parquet::arrow::WriteTable(*table, pool_.get(), sink, table->num_rows(), + writer_properties, arrow_properties.build()) + .ok()); + auto buffer = sink->Finish().ValueOrDie(); + for (bool passthrough : {false, true}) { + SCOPED_TRACE(passthrough); + ASSERT_OK_AND_ASSIGN( + std::unique_ptr reader, + ParquetFileBatchReader::Create( + std::make_shared(buffer), + {{PARQUET_READ_ENABLE_DICTIONARY_PASSTHROUGH, passthrough ? "true" : "false"}}, + /*batch_size=*/8, /*file_metadata=*/nullptr, /*storage_read_bytes=*/nullptr, + pool_, /*hints=*/std::nullopt)); + ASSERT_OK_AND_ASSIGN(auto batch, reader->NextBatch()); + ASSERT_FALSE(BatchReader::IsEofBatch(batch)); + auto array = arrow::ImportArray(batch.first.get(), batch.second.get()).ValueOrDie(); + auto result = checked_pointer_cast(array)->field(0); + ASSERT_EQ(store_schema ? arrow::Type::DICTIONARY : arrow::Type::BINARY, + result->type_id()); + if (store_schema) { + auto dict = checked_pointer_cast(result); + ASSERT_EQ(arrow::Type::BINARY, dict->dictionary()->type_id()); + } + reader->Close(); + } + } +} + +TEST_F(ParquetFileBatchReaderTest, TestColumnarAccessWithBinaryDictionary) { + const std::vector values = {std::string("\x00\xff\x80", 3), "", + std::string("a\0b", 3)}; + arrow::BinaryBuilder builder; + for (const auto& value : values) { + ASSERT_TRUE(builder.Append(value).ok()); + } + std::shared_ptr dictionary; + ASSERT_TRUE(builder.Finish(&dictionary).ok()); + auto indices = + arrow::ipc::internal::json::ArrayFromJSON(arrow::int32(), "[0, 1, null, 2, 0, 1]") + .ValueOrDie(); + auto dict_array = arrow::DictionaryArray::FromArrays(indices, dictionary).ValueOrDie(); + auto schema = arrow::schema({arrow::field("payload", dict_array->type())}); + auto table = arrow::Table::Make(schema, {dict_array}); + auto sink = arrow::io::BufferOutputStream::Create().ValueOrDie(); + auto writer_properties = ::parquet::WriterProperties::Builder().enable_dictionary()->build(); + // Preserve the Arrow dictionary type so the reader returns dictionary-encoded binary. + auto arrow_properties = ::parquet::ArrowWriterProperties::Builder().store_schema()->build(); + ASSERT_TRUE(::parquet::arrow::WriteTable(*table, pool_.get(), sink, /*chunk_size=*/3, + writer_properties, arrow_properties) + .ok()); + auto buffer = sink->Finish().ValueOrDie(); + auto reader = PrepareParquetFileBatchReader(std::make_unique(buffer), + /*options=*/{}, schema, /*predicate=*/nullptr, + /*selection_bitmap=*/std::nullopt, + /*batch_size=*/2); + ASSERT_TRUE(reader); + + auto pool = GetDefaultPool(); + const std::vector expected_indices = {0, 1, -1, 2, 0, 1}; + int32_t row_offset = 0; + while (true) { + ASSERT_OK_AND_ASSIGN(auto batch, reader->NextBatch()); + if (BatchReader::IsEofBatch(batch)) { + break; + } + auto array = arrow::ImportArray(batch.first.get(), batch.second.get()).ValueOrDie(); + auto struct_array = std::dynamic_pointer_cast(array); + ASSERT_TRUE(struct_array); + auto payload = std::dynamic_pointer_cast(struct_array->field(0)); + ASSERT_TRUE(payload); + ASSERT_EQ(arrow::Type::BINARY, payload->dictionary()->type_id()); + auto ctx = std::make_shared(struct_array->fields(), pool); + ColumnarArray column(payload.get(), pool, /*offset=*/0, payload->length()); + for (int32_t pos = 0; pos < struct_array->length(); ++pos) { + ASSERT_LT(row_offset, expected_indices.size()); + int32_t index = expected_indices[row_offset++]; + ColumnarRow row(struct_array->fields(), pool, pos); + ColumnarRowRef row_ref(ctx, pos); + ASSERT_EQ(index == -1, row.IsNullAt(0)); + ASSERT_EQ(index == -1, row_ref.IsNullAt(0)); + ASSERT_EQ(index == -1, column.IsNullAt(pos)); + if (index == -1) { + continue; + } + ASSERT_EQ(values[index], row.GetStringView(0)); + ASSERT_EQ(Bytes(values[index], pool.get()), *row.GetBinary(0)); + ASSERT_EQ(values[index], row_ref.GetStringView(0)); + ASSERT_EQ(Bytes(values[index], pool.get()), *row_ref.GetBinary(0)); + ASSERT_EQ(values[index], column.GetStringView(pos)); + ASSERT_EQ(Bytes(values[index], pool.get()), *column.GetBinary(pos)); + } + } + ASSERT_EQ(expected_indices.size(), row_offset); + reader->Close(); +} + TEST_F(ParquetFileBatchReaderTest, TestNestedStructChildProjectionRecall) { auto f0 = arrow::field("f0", arrow::int32()); auto f1 = arrow::field( @@ -1983,11 +2101,8 @@ TEST_F(ParquetFileBatchReaderTest, TestDictionaryPassthrough) { } TEST_F(ParquetFileBatchReaderTest, TestDictionaryPassthroughSkipsBinaryColumn) { - // Parquet stores STRING and BINARY in the same BYTE_ARRAY leaf and dictionary-encodes both, so - // the gate has to exclude BINARY by logical type. It does, because nothing downstream can read - // `dictionary(int32, binary)`: ColumnarUtils::GetView() asserts on it and returns an empty view - // in a release build, and LiteralConverter rejects it. `f8` is the control - same physical - // type, same pages, and it is forwarded - so this fails if the exclusion is ever widened back. + // Keep the passthrough policy limited to STRING even though consumers also support binary + // dictionaries. Both columns use BYTE_ARRAY and dictionary pages; f8 is the STRING control. WriteArray(file_path_, struct_array_, schema_, /*write_batch_size=*/struct_array_->length(), /*enable_dictionary=*/true, /*max_row_group_length=*/struct_array_->length()); diff --git a/src/paimon/testing/utils/CMakeLists.txt b/src/paimon/testing/utils/CMakeLists.txt index 25f99d3f..29ad1c77 100644 --- a/src/paimon/testing/utils/CMakeLists.txt +++ b/src/paimon/testing/utils/CMakeLists.txt @@ -45,6 +45,7 @@ if(PAIMON_BUILD_TESTS) add_paimon_test(test_utils_test SOURCES data_generator_test.cpp + dict_array_converter_test.cpp STATIC_LINK_LIBS paimon_shared ${PAIMON_LOCAL_FILE_SYSTEM_SHARED_LINK_LIBS} diff --git a/src/paimon/testing/utils/dict_array_converter.h b/src/paimon/testing/utils/dict_array_converter.h index b2267923..286a5ef4 100644 --- a/src/paimon/testing/utils/dict_array_converter.h +++ b/src/paimon/testing/utils/dict_array_converter.h @@ -23,6 +23,7 @@ #include "arrow/api.h" #include "paimon/common/utils/arrow/status_utils.h" #include "paimon/common/utils/checked_cast.h" +#include "paimon/core/casting/casting_utils.h" #include "paimon/result.h" namespace paimon::test { @@ -31,8 +32,8 @@ class DictArrayConverter { DictArrayConverter() = delete; ~DictArrayConverter() = delete; - // Decode dictionary string arrays to plain StringArray so test comparisons are stable across - // Arrow dictionary index types and string/large_string dictionary values. + // Decode dictionaries recursively, preserving binary values and normalizing string offsets + // so test comparisons are stable across readers. static Result> ConvertDictArray( const std::shared_ptr& array, arrow::MemoryPool* pool) { arrow::Type::type kind = array->type_id(); @@ -87,48 +88,12 @@ class DictArrayConverter { item_array, map_array->null_bitmap(), map_array->null_count(), map_array->offset()); } - case arrow::Type::type::DICTIONARY: { - auto dict_array = checked_pointer_cast(array); - auto dict_type = checked_pointer_cast(dict_array->type()); - auto value_type_id = dict_type->value_type()->id(); - if (value_type_id == arrow::Type::type::STRING) { - return ConvertDictionaryArrayToStringArray(dict_array, - pool); - } else if (value_type_id == arrow::Type::type::LARGE_STRING) { - return ConvertDictionaryArrayToStringArray(dict_array, - pool); - } else { - return Status::Invalid( - "only support STRING or LARGE_STRING value type for DictionaryArray"); - } - } + case arrow::Type::type::DICTIONARY: + return CastingUtils::DecodeDictionary(array, pool); default: { return array; } } } - - private: - template - static Result> ConvertDictionaryArrayToStringArray( - const std::shared_ptr& dict_array, arrow::MemoryPool* pool) { - auto dictionary = std::dynamic_pointer_cast(dict_array->dictionary()); - if (!dictionary) { - return Status::Invalid("dictionary value array type does not match dictionary type"); - } - - arrow::StringBuilder string_builder(pool); - for (int64_t i = 0; i < dict_array->length(); ++i) { - if (dict_array->IsNull(i)) { - PAIMON_RETURN_NOT_OK_FROM_ARROW(string_builder.AppendNull()); - } else { - PAIMON_RETURN_NOT_OK_FROM_ARROW( - string_builder.Append(dictionary->GetString(dict_array->GetValueIndex(i)))); - } - } - std::shared_ptr string_array; - PAIMON_RETURN_NOT_OK_FROM_ARROW(string_builder.Finish(&string_array)); - return string_array; - } }; } // namespace paimon::test diff --git a/src/paimon/testing/utils/dict_array_converter_test.cpp b/src/paimon/testing/utils/dict_array_converter_test.cpp new file mode 100644 index 00000000..e53b01c4 --- /dev/null +++ b/src/paimon/testing/utils/dict_array_converter_test.cpp @@ -0,0 +1,80 @@ +/* + * 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/testing/utils/dict_array_converter.h" + +#include "arrow/ipc/json_simple.h" +#include "gtest/gtest.h" +#include "paimon/testing/utils/testharness.h" + +namespace paimon::test { +TEST(DictArrayConverterTest, TestNestedBinaryDictionary) { + const std::string bytes("a\0\xff\x80", 4); + arrow::BinaryBuilder builder; + ASSERT_TRUE(builder.Append(bytes).ok()); + ASSERT_TRUE(builder.Append("").ok()); + ASSERT_TRUE(builder.AppendNull().ok()); + std::shared_ptr binary_dictionary; + ASSERT_TRUE(builder.Finish(&binary_dictionary).ok()); + for (const auto& value_type : {arrow::binary(), arrow::large_binary()}) { + SCOPED_TRACE(value_type->ToString()); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr dictionary, + CastingUtils::Cast(binary_dictionary, value_type, arrow::compute::CastOptions::Safe(), + arrow::default_memory_pool())); + auto indices = + arrow::ipc::internal::json::ArrayFromJSON(arrow::int16(), "[1, 0, null, 1, 2, 0]") + .ValueOrDie(); + auto encoded = + arrow::DictionaryArray::FromArrays(indices, dictionary).ValueOrDie()->Slice(1); + auto offsets = + arrow::ipc::internal::json::ArrayFromJSON(arrow::int32(), "[0, 2, 2, 5]").ValueOrDie(); + auto list = arrow::ListArray::FromArrays(*offsets, *encoded).ValueOrDie(); + auto keys = + arrow::ipc::internal::json::ArrayFromJSON(arrow::utf8(), R"(["a", "b", "c", "d", "e"])") + .ValueOrDie(); + auto map = arrow::MapArray::FromArrays(offsets, keys, encoded).ValueOrDie(); + auto array = + arrow::StructArray::Make({list, map}, std::vector({"list", "map"})) + .ValueOrDie(); + ASSERT_OK_AND_ASSIGN( + std::shared_ptr result, + DictArrayConverter::ConvertDictArray(array, arrow::default_memory_pool())); + ASSERT_TRUE(result->ValidateFull().ok()); + auto decoded_struct = checked_pointer_cast(result); + auto decoded_list = checked_pointer_cast(decoded_struct->field(0)); + auto decoded_map = checked_pointer_cast(decoded_struct->field(1)); + ASSERT_TRUE(decoded_list->value_type()->Equals(value_type)); + ASSERT_TRUE(decoded_map->items()->type()->Equals(value_type)); + ASSERT_TRUE(decoded_list->values()->Equals(decoded_map->items())); + auto values = decoded_list->values(); + ASSERT_TRUE(values->IsNull(1)); + ASSERT_TRUE(values->IsNull(3)); + auto check_values = [&](const auto& typed) { + ASSERT_EQ(bytes, typed->GetView(0)); + ASSERT_EQ("", typed->GetView(2)); + ASSERT_EQ(bytes, typed->GetView(4)); + }; + if (value_type->id() == arrow::Type::BINARY) { + check_values(checked_pointer_cast(values)); + } else { + check_values(checked_pointer_cast(values)); + } + } +} +} // namespace paimon::test