Skip to content
Open
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
7 changes: 4 additions & 3 deletions docs/source/user_guide/compaction.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
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 @@ -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<arrow::StringArray*>(typed_array->dictionary().get());
checked_cast<arrow::BinaryArray*>(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<arrow::LargeStringArray*>(typed_array->dictionary().get());
checked_cast<arrow::LargeBinaryArray*>(typed_array->dictionary().get());
return dictionary->GetView(dict_index);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Here we call the dictionary’s GetView() directly. If the index itself is valid but the referenced dictionary slot is null, DictionaryArray::IsNull() still returns false. ColumnarRow, ColumnarArray, and ColumnarRowRef’s IsNullAt() only check the outer index validity, so the null slot will later be read as a zero-length view, which is indistinguishable from a valid empty byte string.

So I’d like to confirm whether there are cases where the dictionary indices are not null, but the dictionary entries themselves are null? Since many other places already take this into account and have tests for it.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

By definition, arrow dictionary can have both index and dictionary slot be null, and both resulting a null value. But the common practice (like in parquet) represents null values as null indices. Considering dropping the second check to have the consistent behavior. Thanks.

Also I believe you meant GetLiteralFromDictionaryArray in iteral_converter.h instead of the codes here :)

}
assert(false);
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 @@ -18,6 +18,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 @@ -53,4 +54,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
23 changes: 16 additions & 7 deletions src/paimon/common/predicate/literal_converter.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -214,17 +214,26 @@ Result<std::vector<Literal>> LiteralConverter::ConvertLiteralsFromArray(const ar
auto* dict_type = checked_cast<arrow::DictionaryType*>(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<arrow::StringArray, arrow::Int32Array>(
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<arrow::BinaryArray, arrow::Int32Array>(
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<arrow::LargeStringArray, arrow::Int64Array>(
dict_array, FieldType::STRING, own_data);
FieldType literal_type = value_type_id == arrow::Type::type::LARGE_STRING
? FieldType::STRING
: FieldType::BINARY;
return GetLiteralFromDictionaryArray<arrow::LargeBinaryArray, arrow::Int64Array>(
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:
Expand Down
15 changes: 5 additions & 10 deletions src/paimon/common/predicate/literal_converter.h
Original file line number Diff line number Diff line change
Expand Up @@ -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<DictArrayType, arrow::StringArray>) {
int32_t length = 0;
const uint8_t* value = dictionary->GetValue(dict_index, &length);
literals.emplace_back(literal_type, reinterpret_cast<const char*>(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<const char*>(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;
Expand Down
53 changes: 53 additions & 0 deletions src/paimon/common/predicate/literal_converter_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -449,6 +449,59 @@ TEST_F(LiteralConverterTest, TestDictType) {
Literal(FieldType::STRING, "foo", 3), Literal(FieldType::STRING)}));
}

TEST_F(LiteralConverterTest, TestDictionaryValueAndIndexTypes) {
const std::vector<std::pair<std::shared_ptr<arrow::DataType>, std::shared_ptr<arrow::DataType>>>
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<Literal> 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<Literal> 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<arrow::Array> 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<Literal> actual,
LiteralConverter::ConvertLiteralsFromArray(*array, /*own_data=*/true));
array.reset();
dictionary.reset();
ASSERT_EQ(actual,
std::vector<Literal>({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.
Expand Down
15 changes: 9 additions & 6 deletions src/paimon/common/predicate/null_false_leaf_binary_function.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -98,12 +98,15 @@ std::optional<FieldType> ArrayFieldType(const arrow::Array& array) {
const std::shared_ptr<arrow::DataType>& type = array.type();
if (type->id() == arrow::Type::DICTIONARY) {
const auto& dict_type = checked_cast<const arrow::DictionaryType&>(*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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -339,6 +339,28 @@ TEST_F(NullFalseLeafBinaryFunctionTest, TestLargeStringDictionary) {
std::vector<char>({0, 1, 0, 0}));
}

TEST_F(NullFalseLeafBinaryFunctionTest, TestBinaryDictionary) {
const std::vector<std::pair<std::shared_ptr<arrow::DataType>, std::shared_ptr<arrow::DataType>>>
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<char>({1, 0, 0, 0, 0, 1}));
ASSERT_EQ(Eval(LessThan::Instance(), literal, array),
std::vector<char>({0, 0, 1, 1, 0, 0}));
ASSERT_EQ(Eval(NotEqual::Instance(), BinaryLiteral(""), array),
std::vector<char>({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
Expand Down
13 changes: 13 additions & 0 deletions src/paimon/core/casting/casting_utils.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,19 @@
#include "paimon/common/utils/checked_cast.h"

namespace paimon {
Result<std::shared_ptr<arrow::Array>> CastingUtils::DecodeDictionary(
const std::shared_ptr<arrow::Array>& array, arrow::MemoryPool* pool) {
if (array->type_id() != arrow::Type::DICTIONARY) {
return array;
}
const auto& dictionary_type = checked_cast<const arrow::DictionaryType&>(*array->type());
std::shared_ptr<arrow::DataType> value_type = dictionary_type.value_type();
if (value_type->id() == arrow::Type::LARGE_STRING) {
value_type = arrow::utf8();
}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could value_type->id() be large binary? If so, should it be cast to binary? Since in paimon-cpp, large binary specifically refers to the blob type.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, the dictionary value type can be large_binary. But as suggested in src/paimon/common/utils/field_type_utils.h, Paimon BLOB is expecting Arrow LARGE_BINARY.
Converting large_binary to binary here would lose BLOB semantics, thus LARGE_BINARY is intentionally preserved instead of being cast. And we have a unit test making sure large binary is not cast.
This is different from casting LARGE_STRING to string.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the explanation. Please add a comment here explaining that, because of BLOB type, we keep the original type for LARGE_BINARY here.

return Cast(array, value_type, arrow::compute::CastOptions::Safe(), pool);
}

Result<std::shared_ptr<arrow::Array>> CastingUtils::Cast(
const std::shared_ptr<arrow::Array>& src_array,
const std::shared_ptr<arrow::DataType>& target_type, const arrow::compute::CastOptions& options,
Expand Down
6 changes: 6 additions & 0 deletions src/paimon/core/casting/casting_utils.h
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,12 @@ class PAIMON_EXPORT CastingUtils {
const std::shared_ptr<arrow::DataType>& 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<std::shared_ptr<arrow::Array>> DecodeDictionary(
const std::shared_ptr<arrow::Array>& array, arrow::MemoryPool* pool);

template <typename SrcScalar, typename SrcDataType, typename TargetScalar,
typename TargetDataType>
static Result<Literal> Cast(const Literal& literal,
Expand Down
53 changes: 53 additions & 0 deletions src/paimon/core/casting/casting_utils_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<arrow::Array> 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<arrow::Array> 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<arrow::Array> 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<arrow::Array> 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<arrow::BinaryArray>(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(
Expand Down
Loading
Loading