Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 6 additions & 4 deletions src/paimon/common/data/columnar/columnar_utils.h
Original file line number Diff line number Diff line change
Expand Up @@ -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<arrow::StringArray*>(
if (value_type_id == arrow::Type::type::STRING ||
value_type_id == arrow::Type::type::BINARY) {
auto dictionary = arrow::internal::checked_cast<arrow::BinaryArray*>(
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<arrow::LargeStringArray*>(
} else if (value_type_id == arrow::Type::type::LARGE_STRING ||
value_type_id == arrow::Type::type::LARGE_BINARY) {
auto dictionary = arrow::internal::checked_cast<arrow::LargeBinaryArray*>(
typed_array->dictionary().get());
assert(dictionary);
return dictionary->GetView(dict_index);
Expand Down
48 changes: 48 additions & 0 deletions src/paimon/common/data/columnar/columnar_utils_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
#include "paimon/common/data/columnar/columnar_utils.h"

#include <string>
#include <vector>

#include "arrow/api.h"
#include "arrow/array/array_dict.h"
Expand Down Expand Up @@ -52,4 +53,51 @@ TEST(ColumnarUtilsTest, TestGetViewAndBytesOfDict) {
ASSERT_EQ("foo", std::string(ColumnarUtils::GetView(dict_array.get(), 4)));
}

template <typename ArrowType>
class ColumnarUtilsBinaryDictionaryTest : public ::testing::Test {};

using BinaryDictionaryTypes = ::testing::Types<arrow::BinaryType, arrow::LargeBinaryType>;
TYPED_TEST_SUITE(ColumnarUtilsBinaryDictionaryTest, BinaryDictionaryTypes);

TYPED_TEST(ColumnarUtilsBinaryDictionaryTest, GetViewAndBytes) {
auto pool = GetDefaultPool();
const std::vector<std::string> values = {std::string("\x00\xff\x80", 3), "",
std::string("a\0b", 3)};
typename arrow::TypeTraits<TypeParam>::BuilderType builder;
ASSERT_TRUE(builder.Append("unused").ok());
for (const auto& value : values) {
ASSERT_TRUE(builder.Append(value).ok());
}
std::shared_ptr<arrow::Array> dictionary;
ASSERT_TRUE(builder.Finish(&dictionary).ok());
dictionary = dictionary->Slice(1);

const std::vector<std::shared_ptr<arrow::DataType>> index_types = {
arrow::int8(), arrow::int16(), arrow::int32(), arrow::int64()};
const std::vector<int32_t> 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<TypeParam>(sliced.get(), pos, pool.get());
ASSERT_EQ(Bytes(values[index], pool.get()), *bytes);
}
}
}
}

} // namespace paimon::test
72 changes: 72 additions & 0 deletions src/paimon/format/parquet/parquet_file_batch_reader_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -565,6 +569,74 @@ TEST_F(ParquetFileBatchReaderTest, TestNextBatchWithDictionary) {
check_result(false);
}

TEST_F(ParquetFileBatchReaderTest, TestColumnarAccessWithBinaryDictionary) {
const std::vector<std::string> 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<arrow::Array> 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<arrow::io::BufferReader>(buffer),
/*options=*/{}, schema, /*predicate=*/nullptr,
/*selection_bitmap=*/std::nullopt,
/*batch_size=*/2);
ASSERT_TRUE(reader);

auto pool = GetDefaultPool();
const std::vector<int32_t> 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<arrow::StructArray>(array);
ASSERT_TRUE(struct_array);
auto payload = std::dynamic_pointer_cast<arrow::DictionaryArray>(struct_array->field(0));
ASSERT_TRUE(payload);
ASSERT_EQ(arrow::Type::BINARY, payload->dictionary()->type_id());
auto ctx = std::make_shared<ColumnarBatchContext>(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(
Expand Down
Loading