Skip to content
Merged
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
4 changes: 3 additions & 1 deletion src/iceberg/avro/avro_data_util.cc
Original file line number Diff line number Diff line change
Expand Up @@ -594,8 +594,10 @@ Status ExtractDatumFromArray(const ::arrow::Array& array, int64_t index,
internal::checked_cast<const ::arrow::Decimal128Array&>(array);
std::string_view decimal_value = decimal_array.GetView(index);
auto& fixed_datum = datum->value<::avro::GenericFixed>();
const auto fixed_size =
static_cast<std::ptrdiff_t>(fixed_datum.schema()->fixedSize());
auto& bytes = fixed_datum.value();
bytes.assign(decimal_value.begin(), decimal_value.end());
bytes.assign(decimal_value.begin(), decimal_value.begin() + fixed_size);
std::ranges::reverse(bytes);
return {};
}
Expand Down
4 changes: 3 additions & 1 deletion src/iceberg/avro/avro_direct_encoder.cc
Original file line number Diff line number Diff line change
Expand Up @@ -195,7 +195,9 @@ Status EncodeArrowToAvro(const ::avro::NodePtr& avro_node, ::avro::Encoder& enco
const auto& decimal_array =
internal::checked_cast<const ::arrow::Decimal128Array&>(array);
std::string_view decimal_value = decimal_array.GetView(row_index);
ctx.bytes_scratch.assign(decimal_value.begin(), decimal_value.end());
const auto fixed_size = static_cast<std::ptrdiff_t>(avro_node->fixedSize());
ctx.bytes_scratch.assign(decimal_value.begin(),
decimal_value.begin() + fixed_size);
// Arrow Decimal128 bytes are in little-endian order, Avro requires big-endian
std::ranges::reverse(ctx.bytes_scratch);
encoder.encodeFixed(ctx.bytes_scratch.data(), ctx.bytes_scratch.size());
Expand Down
13 changes: 10 additions & 3 deletions src/iceberg/manifest/manifest_adapter.cc
Original file line number Diff line number Diff line change
Expand Up @@ -227,10 +227,17 @@ Status ManifestEntryAdapter::AppendPartitionValues(
ICEBERG_RETURN_UNEXPECTED(
AppendField(child_array, std::get<int64_t>(partition_value.value())));
break;
case TypeId::kDecimal:
ICEBERG_RETURN_UNEXPECTED(AppendField(
child_array, std::get<Decimal>(partition_value.value()).ToBytes()));
case TypeId::kDecimal: {
const auto& decimal_type =
internal::checked_cast<const DecimalType&>(*partition_field.type());
ArrowDecimal decimal;
ArrowDecimalInit(&decimal, 128, decimal_type.precision(), decimal_type.scale());
ArrowDecimalSetBytes(&decimal,
std::get<Decimal>(partition_value.value()).ToBytes().data());
ICEBERG_NANOARROW_RETURN_UNEXPECTED(
ArrowArrayAppendDecimal(child_array, &decimal));
break;
}
case TypeId::kUuid:
ICEBERG_RETURN_UNEXPECTED(
AppendField(child_array, std::get<Uuid>(partition_value.value()).bytes()));
Expand Down
10 changes: 10 additions & 0 deletions src/iceberg/manifest/manifest_reader.cc
Original file line number Diff line number Diff line change
Expand Up @@ -424,6 +424,16 @@ Status ParsePartitionValues(ArrowArrayView* view, int64_t row_idx,
partition.AddValue(Literal::Binary(std::vector<uint8_t>(
buf_value.data.as_char, buf_value.data.as_char + buf_value.size_bytes)));
} break;
case ArrowType::NANOARROW_TYPE_DECIMAL128: {
const auto& decimal_type = internal::checked_cast<const DecimalType&>(*field_type);
ArrowDecimal decimal;
ArrowDecimalInit(&decimal, 128, decimal_type.precision(), decimal_type.scale());
ArrowArrayViewGetDecimalUnsafe(view, row_idx, &decimal);
const Decimal value(static_cast<int64_t>(decimal.words[decimal.high_word_index]),
decimal.words[decimal.low_word_index]);
partition.AddValue(Literal::Decimal(value.value(), decimal_type.precision(),
decimal_type.scale()));
} break;
default:
return InvalidManifest("Unsupported type {} for partition values",
ArrowTypeString(view->storage_type));
Expand Down
1 change: 1 addition & 0 deletions src/iceberg/test/avro_data_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -920,6 +920,7 @@ const std::vector<ExtractDatumParam> kExtractDatumTestCases = {
const auto& fixed = record.fieldAt(0).value<::avro::GenericFixed>();

const auto& bytes = fixed.value();
EXPECT_EQ(bytes.size(), 5);
auto decimal =
::arrow::Decimal128::FromBigEndian(
reinterpret_cast<const uint8_t*>(bytes.data()), bytes.size())
Expand Down
65 changes: 65 additions & 0 deletions src/iceberg/test/avro_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -882,6 +882,71 @@ TEST_P(AvroWriterTest, WritePrimitiveTypes) {
VerifyWrittenData(test_data);
}

TEST_P(AvroWriterTest, WriteDecimalTypes) {
// Precisions at the boundaries where the Avro fixed size changes.
struct DecimalCase {
int32_t precision;
int32_t scale;
size_t fixed_size;
};
const std::vector<DecimalCase> cases = {
{.precision = 1, .scale = 0, .fixed_size = 1},
{.precision = 2, .scale = 0, .fixed_size = 1},
{.precision = 3, .scale = 1, .fixed_size = 2},
{.precision = 9, .scale = 2, .fixed_size = 4},
{.precision = 10, .scale = 2, .fixed_size = 5},
{.precision = 18, .scale = 4, .fixed_size = 8},
{.precision = 19, .scale = 0, .fixed_size = 9},
{.precision = 28, .scale = 6, .fixed_size = 12},
{.precision = 38, .scale = 10, .fixed_size = 16}};

std::vector<SchemaField> fields;
for (size_t i = 0; i < cases.size(); ++i) {
fields.push_back(
SchemaField::MakeRequired(static_cast<int32_t>(i + 1), "d" + std::to_string(i),
decimal(cases[i].precision, cases[i].scale)));
}
auto schema = std::make_shared<iceberg::Schema>(std::move(fields));

// Max, min, zero and one-ulp values for every precision. Arrow's JSON parser
// requires the string to carry exactly the type's scale.
auto with_scale = [](std::string digits, int32_t scale) {
if (scale > 0) digits.insert(digits.size() - scale, ".");
return digits;
};
auto max_value = [&](const DecimalCase& c) {
return with_scale(std::string(c.precision, '9'), c.scale);
};
auto zero = [&](const DecimalCase& c) {
return with_scale(std::string(c.scale + 1, '0'), c.scale);
};
auto one_ulp = [&](const DecimalCase& c) {
return with_scale(std::string(c.scale, '0') + "1", c.scale);
};
auto row = [&](auto value_of) {
std::string out = "[";
for (size_t i = 0; i < cases.size(); ++i) {
out += (i ? ", \"" : "\"") + value_of(cases[i]) + "\"";
}
return out + "]";
};
std::string test_data =
"[" + row(max_value) + ", " +
row([&](const DecimalCase& c) { return "-" + max_value(c); }) + ", " + row(zero) +
", " + row([&](const DecimalCase& c) { return "-" + one_ulp(c); }) + "]";

WriteAvroFile(schema, test_data);

auto root = PhysicalAvroSchema().root();
ASSERT_EQ(root->leaves(), cases.size());
for (size_t i = 0; i < cases.size(); ++i) {
EXPECT_EQ(root->leafAt(i)->fixedSize(), cases[i].fixed_size)
<< "decimal(" << cases[i].precision << ", " << cases[i].scale << ")";
}

VerifyWrittenData(test_data);
}

TEST_P(AvroWriterTest, WriteUuidType) {
auto schema = std::make_shared<iceberg::Schema>(std::vector<SchemaField>{
SchemaField::MakeRequired(1, "uuid_col", iceberg::uuid())});
Expand Down
29 changes: 29 additions & 0 deletions src/iceberg/test/manifest_writer_versions_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -471,6 +471,35 @@ TEST_F(ManifestWriterVersionsTest, TestV2Write) {
DataFile::Content::kData);
}

TEST_F(ManifestWriterVersionsTest, TestV2WriteDecimalPartitionValues) {
constexpr int128_t kMax38 = [] {
int128_t value = 1;
for (int i = 0; i < 38; ++i) {
value *= 10;
}
return value - 1;
}();
schema_ = std::make_shared<Schema>(
std::vector<SchemaField>{SchemaField::MakeRequired(1, "id", int64()),
SchemaField::MakeRequired(2, "timestamp", timestamp_tz()),
SchemaField::MakeRequired(3, "category", string()),
SchemaField::MakeRequired(4, "data", string()),
SchemaField::MakeRequired(5, "double", float64()),
SchemaField::MakeOptional(6, "narrow", decimal(10, 2)),
SchemaField::MakeOptional(7, "wide", decimal(38, 10))});
spec_ =
PartitionSpec::Make(0, {PartitionField(6, 1000, "narrow", Transform::Identity()),
PartitionField(7, 1001, "wide", Transform::Identity())})
.value();
data_file_->partition = PartitionValues(
{Literal::Decimal(-1234, 10, 2), Literal::Decimal(-kMax38, 38, 10)});

auto entries = ReadManifest(WriteManifest(/*format_version=*/2, {data_file_}));

ASSERT_EQ(entries.size(), 1);
EXPECT_EQ(entries[0].data_file->partition, data_file_->partition);
}

TEST_F(ManifestWriterVersionsTest, TestV2WriteWithInheritance) {
auto manifests =
WriteAndReadManifests({WriteManifest(/*format_version=*/2, {data_file_})}, 2);
Expand Down
Loading