Skip to content
Open
139 changes: 139 additions & 0 deletions src/iceberg/avro/avro_data_util.cc
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,15 @@
* under the License.
*/

#include <span>

#include <arrow/array/builder_binary.h>
#include <arrow/array/builder_decimal.h>
#include <arrow/array/builder_nested.h>
#include <arrow/array/builder_primitive.h>
#include <arrow/extension_type.h>
#include <arrow/json/from_string.h>
#include <arrow/scalar.h>
#include <arrow/type.h>
#include <arrow/util/decimal.h>
#include <avro/Generic.hh>
Expand All @@ -31,6 +34,7 @@
#include <avro/Types.hh>

#include "iceberg/arrow/arrow_status_internal.h"
#include "iceberg/arrow/literal_util_internal.h"
#include "iceberg/avro/avro_data_util_internal.h"
#include "iceberg/avro/avro_schema_util_internal.h"
#include "iceberg/metadata_columns.h"
Expand Down Expand Up @@ -88,6 +92,8 @@ Status AppendStructToBuilder(const ::avro::NodePtr& avro_node,
metadata_context, field_builder));
} else if (field_projection.kind == FieldProjection::Kind::kNull) {
ICEBERG_ARROW_RETURN_NOT_OK(field_builder->AppendNull());
} else if (field_projection.kind == FieldProjection::Kind::kDefault) {
ICEBERG_RETURN_UNEXPECTED(AppendDefaultToBuilder(field_projection, field_builder));
} else if (field_projection.kind == FieldProjection::Kind::kMetadata) {
int32_t field_id = expected_field.field_id();
if (field_id == MetadataColumns::kFilePathColumnId) {
Expand Down Expand Up @@ -466,6 +472,10 @@ Status AppendFieldToBuilder(const ::avro::NodePtr& avro_node,
return {};
}

if (projection.kind == FieldProjection::Kind::kDefault) {
return AppendDefaultToBuilder(projection, array_builder);
}

const bool is_row_lineage =
MetadataColumns::IsRowLineageColumn(projected_field.field_id());

Expand Down Expand Up @@ -497,6 +507,135 @@ Status AppendFieldToBuilder(const ::avro::NodePtr& avro_node,

} // namespace

namespace {

Result<std::shared_ptr<::arrow::Scalar>> MakeDefaultScalar(
const Literal& literal, const std::shared_ptr<::arrow::DataType>& builder_type) {
// The builder's own memory pool is not exposed, so the small scalar buffer uses the
// default pool.
ICEBERG_ASSIGN_OR_RAISE(std::shared_ptr<::arrow::Scalar> scalar,
arrow::ToArrowScalar(literal, ::arrow::default_memory_pool()));

Comment thread
huan233usc marked this conversation as resolved.
// For an extension builder (e.g. `arrow.uuid`) target its storage type: ToArrowScalar
// yields the storage scalar (fixed_size_binary(16) for uuid) and Scalar::CastTo has no
// kernel that targets an extension type. This mirrors MakeDefaultArray's extension
// handling.
std::shared_ptr<::arrow::DataType> target_type = builder_type;
if (target_type->id() == ::arrow::Type::EXTENSION) {
target_type = internal::checked_cast<const ::arrow::ExtensionType&>(*target_type)
.storage_type();
}

if (!scalar->type->Equals(*target_type)) {
ICEBERG_ARROW_ASSIGN_OR_RETURN(scalar, scalar->CastTo(target_type));
}
return scalar;
}

Status PrepareStructDefaultScalars(std::span<FieldProjection> projections,
::arrow::ArrayBuilder* builder);

// Recurse into whatever nested builder this projection describes, so a default cached for
// a struct field is found at any nesting depth (e.g. `list<list<struct<...>>>`) instead
// of only when the collection's child is an immediate struct.
Status PrepareNestedDefaultScalars(FieldProjection& projection,
::arrow::ArrayBuilder* builder) {
if (projection.kind != FieldProjection::Kind::kProjected ||
projection.children.empty()) {
return {};
}

switch (builder->type()->id()) {
case ::arrow::Type::STRUCT:
return PrepareStructDefaultScalars(projection.children, builder);
case ::arrow::Type::LIST: {
// List projections store a single child for the element.
auto* list_builder = internal::checked_cast<::arrow::ListBuilder*>(builder);
return PrepareNestedDefaultScalars(projection.children[0],
list_builder->value_builder());
}
case ::arrow::Type::LARGE_LIST: {
auto* list_builder = internal::checked_cast<::arrow::LargeListBuilder*>(builder);
return PrepareNestedDefaultScalars(projection.children[0],
list_builder->value_builder());
}
case ::arrow::Type::MAP: {
auto* map_builder = internal::checked_cast<::arrow::MapBuilder*>(builder);
if (projection.children.size() >= 1) {
ICEBERG_RETURN_UNEXPECTED(PrepareNestedDefaultScalars(
projection.children[0], map_builder->key_builder()));
}
if (projection.children.size() >= 2) {
ICEBERG_RETURN_UNEXPECTED(PrepareNestedDefaultScalars(
projection.children[1], map_builder->item_builder()));
}
return {};
}
default:
return {};
}
}

Status PrepareStructDefaultScalars(std::span<FieldProjection> projections,
::arrow::ArrayBuilder* builder) {
auto* struct_builder = internal::checked_cast<::arrow::StructBuilder*>(builder);
if (static_cast<size_t>(struct_builder->num_fields()) != projections.size()) {
return InvalidArgument(
"Inconsistent number of struct builder fields ({}) and projections ({})",
struct_builder->num_fields(), projections.size());
}

for (size_t i = 0; i < projections.size(); ++i) {
auto& field_projection = projections[i];
auto* field_builder = struct_builder->field_builder(static_cast<int>(i));

if (field_projection.kind == FieldProjection::Kind::kDefault) {
if (field_projection.attributes != nullptr) {
continue;
}
auto attrs = std::make_shared<AvroExtraAttributes>();
ICEBERG_ASSIGN_OR_RAISE(attrs->default_scalar,
MakeDefaultScalar(std::get<Literal>(field_projection.from),
field_builder->type()));
field_projection.attributes = std::move(attrs);
continue;
}

ICEBERG_RETURN_UNEXPECTED(
PrepareNestedDefaultScalars(field_projection, field_builder));
}
return {};
}

} // namespace

Status PrepareDefaultScalars(SchemaProjection& projection,
::arrow::ArrayBuilder* root_builder) {
return PrepareStructDefaultScalars(projection.fields, root_builder);
}

Status AppendDefaultToBuilder(const Literal& literal, ::arrow::ArrayBuilder* builder) {
ICEBERG_ASSIGN_OR_RAISE(std::shared_ptr<::arrow::Scalar> scalar,
MakeDefaultScalar(literal, builder->type()));
ICEBERG_ARROW_RETURN_NOT_OK(builder->AppendScalar(*scalar));
return {};
}

Status AppendDefaultToBuilder(const FieldProjection& projection,
::arrow::ArrayBuilder* builder) {
// Avro projections carry a single attributes type, so once one is attached it is an
// AvroExtraAttributes; use checked_cast instead of a per-row dynamic_cast.
if (projection.attributes != nullptr) {
const auto& attrs =
internal::checked_cast<const AvroExtraAttributes&>(*projection.attributes);
if (attrs.default_scalar != nullptr) {
ICEBERG_ARROW_RETURN_NOT_OK(builder->AppendScalar(*attrs.default_scalar));
return {};
}
}
return AppendDefaultToBuilder(std::get<Literal>(projection.from), builder);
}

Status AppendDatumToBuilder(const ::avro::NodePtr& avro_node,
const ::avro::GenericDatum& avro_datum,
const SchemaProjection& projection,
Expand Down
37 changes: 37 additions & 0 deletions src/iceberg/avro/avro_data_util_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -19,14 +19,51 @@

#pragma once

#include <memory>

#include <arrow/array/builder_base.h>
#include <arrow/scalar.h>
#include <avro/GenericDatum.hh>

#include "iceberg/arrow/metadata_column_util_internal.h"
#include "iceberg/expression/literal.h"
#include "iceberg/schema_util.h"

namespace iceberg::avro {

/// \brief Avro-specific per-field projection attributes.
///
/// `FieldProjection` has a single attributes slot, so all Avro-side attributes live in
/// one container (mirroring `ParquetExtraAttributes`) rather than separate subclasses
/// that could not coexist. `default_scalar` is the Arrow scalar for a `kDefault` field,
/// materialized once (see `PrepareDefaultScalars`) so row-by-row decode only needs
/// `AppendScalar` instead of repeating `ToArrowScalar` / `CastTo` per row.
struct AvroExtraAttributes : FieldProjection::ExtraAttributes {
std::shared_ptr<::arrow::Scalar> default_scalar;
};

/// \brief Precompute cast Arrow scalars for every `kDefault` field under `projection`.
///
/// Walks `root_builder` in lockstep with the projection so each default is cast to the
/// builder's Arrow type once per scan. Safe to call repeatedly; existing
/// `AvroExtraAttributes` entries are left unchanged.
Status PrepareDefaultScalars(SchemaProjection& projection,
::arrow::ArrayBuilder* root_builder);

/// \brief Append a literal once to `builder` while decoding Avro row-by-row.
///
/// Used to materialize `FieldProjection::Kind::kDefault`. Shares `ToArrowScalar` with
/// Parquet's batch path (`MakeDefaultArray`); the append shape stays Avro-local because
/// Avro builds Arrow arrays via per-row `ArrayBuilder`s rather than whole-column arrays.
/// Prefer the `FieldProjection` overload after `PrepareDefaultScalars` so the scalar is
/// reused across rows.
Status AppendDefaultToBuilder(const Literal& literal, ::arrow::ArrayBuilder* builder);

/// \brief Append a `kDefault` projection, reusing a scalar cached on
/// `projection.attributes` when present.
Status AppendDefaultToBuilder(const FieldProjection& projection,
::arrow::ArrayBuilder* builder);

/// \brief Append an Avro datum to an Arrow array builder.
///
/// This function handles schema evolution by using the provided projection to map
Expand Down
3 changes: 3 additions & 0 deletions src/iceberg/avro/avro_direct_decoder.cc
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
#include <avro/Types.hh>

#include "iceberg/arrow/arrow_status_internal.h"
#include "iceberg/avro/avro_data_util_internal.h"
#include "iceberg/avro/avro_direct_decoder_internal.h"
#include "iceberg/avro/avro_schema_util_internal.h"
#include "iceberg/metadata_columns.h"
Expand Down Expand Up @@ -209,6 +210,8 @@ Status DecodeStructToBuilder(const ::avro::NodePtr& avro_node, ::avro::Decoder&
auto* field_builder = struct_builder->field_builder(static_cast<int>(proj_idx));
if (field_projection.kind == FieldProjection::Kind::kNull) {
ICEBERG_ARROW_RETURN_NOT_OK(field_builder->AppendNull());
} else if (field_projection.kind == FieldProjection::Kind::kDefault) {
ICEBERG_RETURN_UNEXPECTED(AppendDefaultToBuilder(field_projection, field_builder));
} else if (field_projection.kind == FieldProjection::Kind::kMetadata) {
int32_t field_id = expected_field.field_id();
if (field_id == MetadataColumns::kFilePathColumnId) {
Expand Down
2 changes: 2 additions & 0 deletions src/iceberg/avro/avro_reader.cc
Original file line number Diff line number Diff line change
Expand Up @@ -434,6 +434,8 @@ class AvroReader::Impl {
builder_result.status().message());
}
context_->builder_ = builder_result.MoveValueUnsafe();
ICEBERG_RETURN_UNEXPECTED(
PrepareDefaultScalars(projection_, context_->builder_.get()));
backend_->InitReadContext(backend_->GetReaderSchema());

return {};
Expand Down
4 changes: 4 additions & 0 deletions src/iceberg/avro/avro_schema_util.cc
Original file line number Diff line number Diff line change
Expand Up @@ -752,6 +752,10 @@ Result<FieldProjection> ProjectStruct(const StructType& struct_type,
iter->second.local_index, prune_source));
} else if (MetadataColumns::IsMetadataColumn(field_id)) {
child_projection.kind = FieldProjection::Kind::kMetadata;
} else if (expected_field.initial_default() != nullptr) {
// Rows written before the field existed assume its `initial-default` value.
child_projection.kind = FieldProjection::Kind::kDefault;
child_projection.from = *expected_field.initial_default();
} else if (expected_field.optional()) {
child_projection.kind = FieldProjection::Kind::kNull;
} else {
Expand Down
Loading
Loading