From b8e96281e6c1ed4b6d15be69c9b41492e039ec8c Mon Sep 17 00:00:00 2001 From: Konrad Piotrowski Date: Wed, 30 Sep 2026 19:33:32 +0200 Subject: [PATCH 1/2] fix: write and read decimal partition values in manifests The partition struct child for a decimal field is a nanoarrow DECIMAL128 array, which ArrowArrayAppendBytes rejects with EINVAL, so any manifest entry carrying a decimal partition value failed to write. Append it as an ArrowDecimal instead. Both avro encoders copied all 16 bytes of a decimal128 into a fixed sized for the precision (5 bytes for decimal(10,2)), corrupting the record for any precision below 38. Emit only the fixed's width, big-endian. The manifest reader had no case for decimal partition values, so a manifest carrying one could not be read back. Parse them into decimal literals. --- src/iceberg/avro/avro_data_util.cc | 4 ++- src/iceberg/avro/avro_direct_encoder.cc | 4 ++- src/iceberg/manifest/manifest_adapter.cc | 13 +++++++-- src/iceberg/manifest/manifest_reader.cc | 10 +++++++ src/iceberg/test/avro_data_test.cc | 1 + .../test/manifest_writer_versions_test.cc | 29 +++++++++++++++++++ 6 files changed, 56 insertions(+), 5 deletions(-) diff --git a/src/iceberg/avro/avro_data_util.cc b/src/iceberg/avro/avro_data_util.cc index b44e4b110..e93bb3f70 100644 --- a/src/iceberg/avro/avro_data_util.cc +++ b/src/iceberg/avro/avro_data_util.cc @@ -594,8 +594,10 @@ Status ExtractDatumFromArray(const ::arrow::Array& array, int64_t index, internal::checked_cast(array); std::string_view decimal_value = decimal_array.GetView(index); auto& fixed_datum = datum->value<::avro::GenericFixed>(); + const auto fixed_size = + static_cast(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 {}; } diff --git a/src/iceberg/avro/avro_direct_encoder.cc b/src/iceberg/avro/avro_direct_encoder.cc index 5dcfd2511..1eb68d436 100644 --- a/src/iceberg/avro/avro_direct_encoder.cc +++ b/src/iceberg/avro/avro_direct_encoder.cc @@ -195,7 +195,9 @@ Status EncodeArrowToAvro(const ::avro::NodePtr& avro_node, ::avro::Encoder& enco const auto& decimal_array = internal::checked_cast(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(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()); diff --git a/src/iceberg/manifest/manifest_adapter.cc b/src/iceberg/manifest/manifest_adapter.cc index 526d31271..7154608ad 100644 --- a/src/iceberg/manifest/manifest_adapter.cc +++ b/src/iceberg/manifest/manifest_adapter.cc @@ -227,10 +227,17 @@ Status ManifestEntryAdapter::AppendPartitionValues( ICEBERG_RETURN_UNEXPECTED( AppendField(child_array, std::get(partition_value.value()))); break; - case TypeId::kDecimal: - ICEBERG_RETURN_UNEXPECTED(AppendField( - child_array, std::get(partition_value.value()).ToBytes())); + case TypeId::kDecimal: { + const auto& decimal_type = + internal::checked_cast(*partition_field.type()); + ArrowDecimal decimal; + ArrowDecimalInit(&decimal, 128, decimal_type.precision(), decimal_type.scale()); + ArrowDecimalSetBytes(&decimal, + std::get(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(partition_value.value()).bytes())); diff --git a/src/iceberg/manifest/manifest_reader.cc b/src/iceberg/manifest/manifest_reader.cc index 8757b5d61..c5b0489af 100644 --- a/src/iceberg/manifest/manifest_reader.cc +++ b/src/iceberg/manifest/manifest_reader.cc @@ -424,6 +424,16 @@ Status ParsePartitionValues(ArrowArrayView* view, int64_t row_idx, partition.AddValue(Literal::Binary(std::vector( 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(*field_type); + ArrowDecimal decimal; + ArrowDecimalInit(&decimal, 128, decimal_type.precision(), decimal_type.scale()); + ArrowArrayViewGetDecimalUnsafe(view, row_idx, &decimal); + const Decimal value(static_cast(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)); diff --git a/src/iceberg/test/avro_data_test.cc b/src/iceberg/test/avro_data_test.cc index 7731f58d3..29a88bd6d 100644 --- a/src/iceberg/test/avro_data_test.cc +++ b/src/iceberg/test/avro_data_test.cc @@ -920,6 +920,7 @@ const std::vector 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(bytes.data()), bytes.size()) diff --git a/src/iceberg/test/manifest_writer_versions_test.cc b/src/iceberg/test/manifest_writer_versions_test.cc index 24669d5b5..2f25e5168 100644 --- a/src/iceberg/test/manifest_writer_versions_test.cc +++ b/src/iceberg/test/manifest_writer_versions_test.cc @@ -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( + std::vector{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); From 1152e1f2fa2ad3a4a31838b8a36188ebfb2f975d Mon Sep 17 00:00:00 2001 From: Konrad Piotrowski Date: Fri, 2 Oct 2026 16:20:43 +0200 Subject: [PATCH 2/2] test(avro): round-trip decimals across fixed-size boundaries Cover precisions 1 through 38 at every point where the Avro fixed width changes, with max, min, zero and one-ulp values, on both the direct encoder and GenericDatum write paths. Also assert the physical fixed size per precision. --- src/iceberg/test/avro_test.cc | 65 +++++++++++++++++++++++++++++++++++ 1 file changed, 65 insertions(+) diff --git a/src/iceberg/test/avro_test.cc b/src/iceberg/test/avro_test.cc index 3ae1696d0..d47056d3f 100644 --- a/src/iceberg/test/avro_test.cc +++ b/src/iceberg/test/avro_test.cc @@ -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 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 fields; + for (size_t i = 0; i < cases.size(); ++i) { + fields.push_back( + SchemaField::MakeRequired(static_cast(i + 1), "d" + std::to_string(i), + decimal(cases[i].precision, cases[i].scale))); + } + auto schema = std::make_shared(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(std::vector{ SchemaField::MakeRequired(1, "uuid_col", iceberg::uuid())});