Skip to content
Draft
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
1 change: 1 addition & 0 deletions src/iceberg/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ set(ICEBERG_SOURCES
geospatial.cc
inspect/history_table.cc
inspect/metadata_table.cc
inspect/position_deletes_table.cc
inspect/snapshots_table.cc
inheritable_metadata.cc
json_serde.cc
Expand Down
6 changes: 6 additions & 0 deletions src/iceberg/arrow_c_data_util.cc
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,12 @@ Status AppendValue(const ArrowSchema& input_schema, const ArrowArray& input_arra

} // namespace

Status AppendArrayValue(const ArrowSchema& input_schema, const ArrowArray& input_array,
const ArrowArrayView& input_view, int64_t row_index,
ArrowArray* output_array) {
return AppendValue(input_schema, input_array, input_view, row_index, output_array);
}

ProjectionContext::ProjectionContext(ProjectionContext&& other) noexcept
: input_schema_(std::exchange(other.input_schema_, nullptr)),
output_schema_(std::exchange(other.output_schema_, nullptr)),
Expand Down
8 changes: 8 additions & 0 deletions src/iceberg/arrow_c_data_util_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,8 @@
#include "iceberg/result.h"
#include "iceberg/type_fwd.h"

struct ArrowArrayView;

namespace iceberg {

/// \brief Cached state for ProjectBatch over one input/output schema pair.
Expand Down Expand Up @@ -242,4 +244,10 @@ ICEBERG_EXPORT Result<ArrowArray> ProjectBatch(ArrowArray* input_batch,
std::span<const int32_t> row_indices,
ProjectionContext& projection);

/// \brief Append one value from an Arrow array into a compatible nanoarrow builder.
ICEBERG_EXPORT Status AppendArrayValue(const ArrowSchema& input_schema,
const ArrowArray& input_array,
const ArrowArrayView& input_view,
int64_t row_index, ArrowArray* output_array);

} // namespace iceberg
10 changes: 10 additions & 0 deletions src/iceberg/arrow_row_builder.cc
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,8 @@
#include "iceberg/nanoarrow_status_internal.h"
#include "iceberg/schema.h"
#include "iceberg/schema_internal.h"
#include "iceberg/type.h"
#include "iceberg/util/decimal.h"

namespace iceberg {

Expand Down Expand Up @@ -137,6 +139,14 @@ Status AppendBytes(ArrowArray* array, std::span<const uint8_t> value) {
return {};
}

Status AppendDecimal(ArrowArray* array, const Decimal& value, const DecimalType& type) {
ArrowDecimal decimal;
ArrowDecimalInit(&decimal, 128, type.precision(), type.scale());
ArrowDecimalSetBytes(&decimal, value.native_endian_bytes());
ICEBERG_NANOARROW_RETURN_UNEXPECTED(ArrowArrayAppendDecimal(array, &decimal));
return {};
}

Status AppendIntList(ArrowArray* array, const std::vector<int32_t>& values) {
if (values.empty()) {
return AppendNull(array);
Expand Down
4 changes: 4 additions & 0 deletions src/iceberg/arrow_row_builder_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,10 @@ ICEBERG_EXPORT Status AppendString(ArrowArray* array, std::string_view value);
/// \brief Append a binary value to a nanoarrow array builder.
ICEBERG_EXPORT Status AppendBytes(ArrowArray* array, std::span<const uint8_t> value);

/// \brief Append an unscaled decimal value to a Decimal128 array.
ICEBERG_EXPORT Status AppendDecimal(ArrowArray* array, const Decimal& value,
const DecimalType& type);

/// \brief Append a list of int32 values to a nanoarrow list array builder.
ICEBERG_EXPORT Status AppendIntList(ArrowArray* array,
const std::vector<int32_t>& values);
Expand Down
1 change: 1 addition & 0 deletions src/iceberg/avro/avro_data_util.cc
Original file line number Diff line number Diff line change
Expand Up @@ -758,6 +758,7 @@ 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>();
decimal_value = decimal_value.substr(0, fixed_datum.schema()->fixedSize());
auto& bytes = fixed_datum.value();
bytes.assign(decimal_value.begin(), decimal_value.end());
std::ranges::reverse(bytes);
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 @@ -194,7 +194,9 @@ Status EncodeArrowToAvro(const ::avro::NodePtr& avro_node, ::avro::Encoder& enco
if (avro_node->logicalType().type() == ::avro::LogicalType::DECIMAL) {
const auto& decimal_array =
internal::checked_cast<const ::arrow::Decimal128Array&>(array);
std::string_view decimal_value = decimal_array.GetView(row_index);
// Avro fixed uses the schema's precision-dependent width, not all 16 Arrow bytes.
std::string_view decimal_value =
decimal_array.GetView(row_index).substr(0, avro_node->fixedSize());
ctx.bytes_scratch.assign(decimal_value.begin(), decimal_value.end());
// Arrow Decimal128 bytes are in little-endian order, Avro requires big-endian
std::ranges::reverse(ctx.bytes_scratch);
Expand Down
4 changes: 4 additions & 0 deletions src/iceberg/deletes/position_delete_index.cc
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,10 @@ int64_t PositionDeleteIndex::Cardinality() const {
return static_cast<int64_t>(bitmap_.Cardinality());
}

void PositionDeleteIndex::ForEach(const std::function<void(int64_t)>& fn) const {
bitmap_.ForEach(fn);
}

void PositionDeleteIndex::Merge(const PositionDeleteIndex& other) {
bitmap_.Or(other.bitmap_);
delete_files_.insert(delete_files_.end(), other.delete_files_.begin(),
Expand Down
4 changes: 4 additions & 0 deletions src/iceberg/deletes/position_delete_index.h
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
/// Index of deleted row positions for a data file.

#include <cstdint>
#include <functional>
#include <memory>
#include <span>
#include <vector>
Expand Down Expand Up @@ -66,6 +67,9 @@ class ICEBERG_EXPORT PositionDeleteIndex {
/// \brief Get the number of deleted positions.
int64_t Cardinality() const;

/// \brief Iterate over deleted positions in ascending order.
void ForEach(const std::function<void(int64_t)>& fn) const;

/// \brief Merge another index into this one.
/// \param other The index to merge (union operation)
void Merge(const PositionDeleteIndex& other);
Expand Down
6 changes: 6 additions & 0 deletions src/iceberg/file_reader.h
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,12 @@ class ICEBERG_EXPORT Reader {
/// \brief Get the schema of the data.
virtual Result<ArrowSchema> Schema() = 0;

/// \brief Check whether a top-level field ID exists in the physical file schema.
/// \return Presence, or NotImplemented for readers without schema inspection.
virtual Result<bool> HasTopLevelField(int32_t field_id) {
return NotImplemented("Physical schema inspection is not supported");
}

/// \brief Get the metadata of the file.
virtual Result<std::unordered_map<std::string, std::string>> Metadata() = 0;
};
Expand Down
1 change: 1 addition & 0 deletions src/iceberg/inspect/metadata_table.h
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ class ICEBERG_EXPORT MetadataTable {
enum class Kind {
kSnapshots,
kHistory,
kPositionDeletes,
};

/// \brief Maximum number of rows emitted in each Arrow batch.
Expand Down
Loading