From 94adfec3edff1bf01e549c8488a72ce4f6318275 Mon Sep 17 00:00:00 2001 From: Wei Zhang Date: Thu, 17 Sep 2026 21:03:22 +0800 Subject: [PATCH] fix(data issue): fix dictionary encoded binary read as empty --- .../common/data/columnar/columnar_utils.h | 10 +-- .../data/columnar/columnar_utils_test.cpp | 48 +++++++++++++ .../parquet_file_batch_reader_test.cpp | 72 +++++++++++++++++++ 3 files changed, 126 insertions(+), 4 deletions(-) diff --git a/src/paimon/common/data/columnar/columnar_utils.h b/src/paimon/common/data/columnar/columnar_utils.h index 6a095e415..ee805f235 100644 --- a/src/paimon/common/data/columnar/columnar_utils.h +++ b/src/paimon/common/data/columnar/columnar_utils.h @@ -87,13 +87,15 @@ class ColumnarUtils { dict_index = indices->Value(pos); } assert(dict_index >= 0); - if (value_type_id == arrow::Type::type::STRING) { - auto dictionary = arrow::internal::checked_cast( + if (value_type_id == arrow::Type::type::STRING || + value_type_id == arrow::Type::type::BINARY) { + auto dictionary = arrow::internal::checked_cast( typed_array->dictionary().get()); assert(dictionary); return dictionary->GetView(dict_index); - } else if (value_type_id == arrow::Type::type::LARGE_STRING) { - auto dictionary = arrow::internal::checked_cast( + } else if (value_type_id == arrow::Type::type::LARGE_STRING || + value_type_id == arrow::Type::type::LARGE_BINARY) { + auto dictionary = arrow::internal::checked_cast( typed_array->dictionary().get()); assert(dictionary); return dictionary->GetView(dict_index); diff --git a/src/paimon/common/data/columnar/columnar_utils_test.cpp b/src/paimon/common/data/columnar/columnar_utils_test.cpp index 17ae8ddc4..9dcfa1fae 100644 --- a/src/paimon/common/data/columnar/columnar_utils_test.cpp +++ b/src/paimon/common/data/columnar/columnar_utils_test.cpp @@ -17,6 +17,7 @@ #include "paimon/common/data/columnar/columnar_utils.h" #include +#include #include "arrow/api.h" #include "arrow/array/array_dict.h" @@ -52,4 +53,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/format/parquet/parquet_file_batch_reader_test.cpp b/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp index f284bfed6..ba7be4154 100644 --- a/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp +++ b/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp @@ -32,9 +32,13 @@ #include "arrow/c/bridge.h" #include "arrow/io/caching.h" #include "arrow/io/interfaces.h" +#include "arrow/io/memory.h" #include "arrow/ipc/api.h" #include "arrow/ipc/json_simple.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/types/data_field.h" #include "paimon/common/utils/arrow/arrow_input_stream_adapter.h" #include "paimon/common/utils/arrow/arrow_utils.h" @@ -565,6 +569,74 @@ TEST_F(ParquetFileBatchReaderTest, TestNextBatchWithDictionary) { check_result(false); } +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(