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
1 change: 1 addition & 0 deletions src/iceberg/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -314,6 +314,7 @@ if(ICEBERG_BUILD_BUNDLE)
avro/avro_stream_internal.cc
parquet/parquet_data_util.cc
parquet/parquet_metrics.cc
parquet/parquet_metrics_row_group_filter.cc
parquet/parquet_reader.cc
parquet/parquet_register.cc
parquet/parquet_schema_util.cc
Expand Down
18 changes: 11 additions & 7 deletions src/iceberg/data/file_scan_task_reader.cc
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@
#include <memory>
#include <optional>
#include <utility>
#include <vector>

#include "iceberg/arrow_c_data_guard_internal.h"
#include "iceberg/arrow_c_data_util_internal.h"
Expand All @@ -42,6 +41,7 @@ namespace {
ReaderOptions MakeReaderOptions(const DataFile& data_file, std::shared_ptr<FileIO> io,
std::shared_ptr<Schema> projection,
std::shared_ptr<Expression> filter,
bool filter_case_sensitive,
std::shared_ptr<NameMapping> name_mapping,
ReaderProperties properties,
std::optional<int64_t> first_row_id,
Expand All @@ -52,6 +52,7 @@ ReaderOptions MakeReaderOptions(const DataFile& data_file, std::shared_ptr<FileI
.io = std::move(io),
.projection = std::move(projection),
.filter = std::move(filter),
.filter_case_sensitive = filter_case_sensitive,
.name_mapping = std::move(name_mapping),
.first_row_id = first_row_id,
.data_sequence_number = data_sequence_number,
Expand Down Expand Up @@ -164,10 +165,13 @@ class FileScanTaskReader::Impl {
"Data file size must not be negative: {}",
data_file->file_size_in_bytes);

auto filter = task.residual_filter();

if (task.delete_files().empty()) {
auto options = MakeReaderOptions(
*data_file, io_, projected_schema_, task.residual_filter(), name_mapping_,
properties_, data_file->first_row_id, data_file->data_sequence_number);
auto options =
MakeReaderOptions(*data_file, io_, projected_schema_, filter,
filter_case_sensitive_, name_mapping_, properties_,
data_file->first_row_id, data_file->data_sequence_number);
ICEBERG_ASSIGN_OR_RAISE(
auto reader, ReaderFactoryRegistry::Open(data_file->file_format, options));
return MakeArrowArrayStream(std::move(reader));
Expand All @@ -192,7 +196,7 @@ class FileScanTaskReader::Impl {
project_batch_function));

auto options = MakeReaderOptions(
*data_file, io_, required_schema, task.residual_filter(), name_mapping_,
*data_file, io_, required_schema, filter, filter_case_sensitive_, name_mapping_,
properties_, data_file->first_row_id, data_file->data_sequence_number);
ICEBERG_ASSIGN_OR_RAISE(auto reader,
ReaderFactoryRegistry::Open(data_file->file_format, options));
Expand All @@ -207,16 +211,16 @@ class FileScanTaskReader::Impl {
Impl(Options options, DeleteFilter::FieldLookup field_lookup,
std::shared_ptr<DeleteCounter> delete_counter)
: io_(std::move(options.io)),
schemas_(std::move(options.schemas)),
projected_schema_(std::move(options.projected_schema)),
filter_case_sensitive_(options.filter_case_sensitive),
name_mapping_(std::move(options.name_mapping)),
properties_(ReaderProperties::FromMap(options.properties)),
field_lookup_(std::move(field_lookup)),
delete_counter_(std::move(delete_counter)) {}

std::shared_ptr<FileIO> io_;
std::vector<std::shared_ptr<Schema>> schemas_;
std::shared_ptr<Schema> projected_schema_;
bool filter_case_sensitive_;
std::shared_ptr<NameMapping> name_mapping_;
ReaderProperties properties_;
DeleteFilter::FieldLookup field_lookup_;
Expand Down
2 changes: 2 additions & 0 deletions src/iceberg/data/file_scan_task_reader.h
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,8 @@ class ICEBERG_DATA_EXPORT FileScanTaskReader {
/// The output schema for the returned ArrowArrayStream. Must be a
/// projection of table_schema.
std::shared_ptr<Schema> projected_schema;
/// Case sensitivity for filter binding. Pass TableScan::is_case_sensitive().
bool filter_case_sensitive = true;
/// Optional name mapping for files written without field IDs.
std::shared_ptr<NameMapping> name_mapping;
/// Format-specific or implementation-specific options for data readers.
Expand Down
51 changes: 35 additions & 16 deletions src/iceberg/expression/binder.cc
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,21 @@

namespace iceberg {

namespace {

Result<std::optional<bool>> CombineResults(const std::optional<bool>& left_result,
const std::optional<bool>& right_result) {
if (!left_result.has_value()) {
return right_result;
}
if (right_result.has_value() && left_result != right_result) {
return InvalidExpression("Found partially bound expression");
}
return left_result;
}

} // namespace

Binder::Binder(const Schema& schema, bool case_sensitive)
: schema_(schema), case_sensitive_(case_sensitive) {}

Expand Down Expand Up @@ -82,40 +97,44 @@ Result<std::shared_ptr<Expression>> Binder::Aggregate(
Result<bool> IsBoundVisitor::IsBound(const std::shared_ptr<Expression>& expr) {
ICEBERG_PRECHECK(expr != nullptr, "Expression cannot be null");
IsBoundVisitor visitor;
return Visit<bool, IsBoundVisitor>(expr, visitor);
ICEBERG_ASSIGN_OR_RAISE(auto is_bound, Visit<std::optional<bool>>(expr, visitor));
return is_bound.value_or(false);
}

Result<bool> IsBoundVisitor::AlwaysTrue() {
return InvalidExpression("IsBoundVisitor does not support AlwaysTrue expression");
}
Result<std::optional<bool>> IsBoundVisitor::AlwaysTrue() { return std::nullopt; }

Result<bool> IsBoundVisitor::AlwaysFalse() {
return InvalidExpression("IsBoundVisitor does not support AlwaysFalse expression");
}
Result<std::optional<bool>> IsBoundVisitor::AlwaysFalse() { return std::nullopt; }

Result<bool> IsBoundVisitor::Not(bool child_result) { return child_result; }
Result<std::optional<bool>> IsBoundVisitor::Not(const std::optional<bool>& child_result) {
return child_result;
}

Result<bool> IsBoundVisitor::And(bool left_result, bool right_result) {
return left_result && right_result;
Result<std::optional<bool>> IsBoundVisitor::And(const std::optional<bool>& left_result,
const std::optional<bool>& right_result) {
return CombineResults(left_result, right_result);
}

Result<bool> IsBoundVisitor::Or(bool left_result, bool right_result) {
return left_result && right_result;
Result<std::optional<bool>> IsBoundVisitor::Or(const std::optional<bool>& left_result,
const std::optional<bool>& right_result) {
return CombineResults(left_result, right_result);
}

Result<bool> IsBoundVisitor::Predicate(const std::shared_ptr<BoundPredicate>& pred) {
Result<std::optional<bool>> IsBoundVisitor::Predicate(
const std::shared_ptr<BoundPredicate>& pred) {
return true;
}

Result<bool> IsBoundVisitor::Predicate(const std::shared_ptr<UnboundPredicate>& pred) {
Result<std::optional<bool>> IsBoundVisitor::Predicate(
const std::shared_ptr<UnboundPredicate>& pred) {
return false;
}

Result<bool> IsBoundVisitor::Aggregate(const std::shared_ptr<BoundAggregate>& aggregate) {
Result<std::optional<bool>> IsBoundVisitor::Aggregate(
const std::shared_ptr<BoundAggregate>& aggregate) {
return true;
}

Result<bool> IsBoundVisitor::Aggregate(
Result<std::optional<bool>> IsBoundVisitor::Aggregate(
const std::shared_ptr<UnboundAggregate>& aggregate) {
return false;
}
Expand Down
27 changes: 17 additions & 10 deletions src/iceberg/expression/binder.h
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
/// Bind an expression to a schema.

#include <functional>
#include <optional>
#include <unordered_set>

#include "iceberg/expression/expression_visitor.h"
Expand Down Expand Up @@ -61,19 +62,25 @@ class ICEBERG_EXPORT Binder : public ExpressionVisitor<std::shared_ptr<Expressio
const bool case_sensitive_;
};

class ICEBERG_EXPORT IsBoundVisitor : public ExpressionVisitor<bool> {
class ICEBERG_EXPORT IsBoundVisitor : public ExpressionVisitor<std::optional<bool>> {
public:
static Result<bool> IsBound(const std::shared_ptr<Expression>& expr);

Result<bool> AlwaysTrue() override;
Result<bool> AlwaysFalse() override;
Result<bool> Not(bool child_result) override;
Result<bool> And(bool left_result, bool right_result) override;
Result<bool> Or(bool left_result, bool right_result) override;
Result<bool> Predicate(const std::shared_ptr<BoundPredicate>& pred) override;
Result<bool> Predicate(const std::shared_ptr<UnboundPredicate>& pred) override;
Result<bool> Aggregate(const std::shared_ptr<BoundAggregate>& aggregate) override;
Result<bool> Aggregate(const std::shared_ptr<UnboundAggregate>& aggregate) override;
Result<std::optional<bool>> AlwaysTrue() override;
Result<std::optional<bool>> AlwaysFalse() override;
Result<std::optional<bool>> Not(const std::optional<bool>& child_result) override;
Result<std::optional<bool>> And(const std::optional<bool>& left_result,
const std::optional<bool>& right_result) override;
Result<std::optional<bool>> Or(const std::optional<bool>& left_result,
const std::optional<bool>& right_result) override;
Result<std::optional<bool>> Predicate(
const std::shared_ptr<BoundPredicate>& pred) override;
Result<std::optional<bool>> Predicate(
const std::shared_ptr<UnboundPredicate>& pred) override;
Result<std::optional<bool>> Aggregate(
const std::shared_ptr<BoundAggregate>& aggregate) override;
Result<std::optional<bool>> Aggregate(
const std::shared_ptr<UnboundAggregate>& aggregate) override;
};

using FieldIdsSetRef = std::reference_wrapper<std::unordered_set<int32_t>>;
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 @@ -80,6 +80,9 @@ class ICEBERG_EXPORT ReaderProperties : public ConfigBase<ReaderProperties> {
/// Only the Parquet reader honors this option; other readers ignore it.
/// Default: false (use 32-bit offset list).
inline static Entry<bool> kArrowUseLargeList{"read.arrow.use-large-list", false};
/// \brief Use footer statistics to prune Parquet row groups.
inline static Entry<bool> kParquetRowGroupFilter{
"read.parquet.row-group-filter.enabled", true};
/// \brief Skip GenericDatum in Avro reader for better performance.
/// When true, decode directly from Avro to Arrow without GenericDatum intermediate.
/// Default: true (skip GenericDatum for better performance).
Expand All @@ -103,10 +106,13 @@ struct ICEBERG_EXPORT ReaderOptions {
/// \brief FileIO instance to open the file.
std::shared_ptr<class FileIO> io;
/// \brief The projection schema to read from the file. This field is required.
/// Include columns referenced by unbound filters; readers bind them to this schema.
std::shared_ptr<class Schema> projection;
/// \brief The filter to apply to the data. Reader implementations may ignore this if
/// the file format does not support filtering.
std::shared_ptr<class Expression> filter;
/// \brief Case sensitivity inherited from the scan when binding filter references.
bool filter_case_sensitive = true;
/// \brief Name mapping for schema evolution compatibility. Used when reading files
/// that may have different field names than the current schema.
std::shared_ptr<class NameMapping> name_mapping;
Expand Down
90 changes: 45 additions & 45 deletions src/iceberg/parquet/parquet_metrics.cc
Original file line number Diff line number Diff line change
Expand Up @@ -191,47 +191,6 @@ bool NeedsBoundTruncation(const PrimitiveType& type) {
return type.type_id() == TypeId::kString || type.type_id() == TypeId::kBinary;
}

Result<Literal> StatsValueToLiteral(const ::parquet::ColumnDescriptor& column,
const PrimitiveType& iceberg_type,
const ::parquet::Statistics& stats, bool is_min) {
switch (column.physical_type()) {
case ::parquet::Type::BOOLEAN:
return TypedStatsLiteral<::parquet::BoolStatistics>(
stats, is_min, [](bool value) { return Literal::Boolean(value); });
case ::parquet::Type::INT32:
return TypedStatsLiteral<::parquet::Int32Statistics>(
stats, is_min,
[&](int32_t value) { return Int32StatsLiteral(value, iceberg_type); });
case ::parquet::Type::INT64:
return TypedStatsLiteral<::parquet::Int64Statistics>(
stats, is_min,
[&](int64_t value) { return Int64StatsLiteral(value, iceberg_type); });
case ::parquet::Type::FLOAT:
return TypedStatsLiteral<::parquet::FloatStatistics>(
stats, is_min,
[&](float value) { return FloatStatsLiteral(value, iceberg_type); });
case ::parquet::Type::DOUBLE:
return TypedStatsLiteral<::parquet::DoubleStatistics>(
stats, is_min, [](double value) { return Literal::Double(value); });
case ::parquet::Type::BYTE_ARRAY:
return TypedStatsLiteral<::parquet::ByteArrayStatistics>(
stats, is_min, [&](const ::parquet::ByteArray& value) {
return BinaryStatsLiteral(BytesFromByteArray(value), iceberg_type);
});
case ::parquet::Type::FIXED_LEN_BYTE_ARRAY:
return TypedStatsLiteral<::parquet::FLBAStatistics>(
stats, is_min, [&](const ::parquet::FixedLenByteArray& value) {
return BinaryStatsLiteral(BytesFromFLBA(value, column.type_length()),
iceberg_type);
});
case ::parquet::Type::INT96:
case ::parquet::Type::UNDEFINED:
return NotSupported("Cannot convert Parquet statistics for physical type {}",
static_cast<int>(column.physical_type()));
}
std::unreachable();
}

/// \brief Collect counts (value count and null count) from footer statistics.
/// \param field_id The Iceberg field ID.
/// \param metadata The Parquet file metadata.
Expand Down Expand Up @@ -288,15 +247,15 @@ Result<std::optional<FieldMetrics>> CollectBounds(
value_count += column_chunk->num_values();

if (stats->HasMinMax()) {
ICEBERG_ASSIGN_OR_RAISE(auto min_value,
StatsValueToLiteral(*column_desc, *iceberg_type, *stats,
ICEBERG_ASSIGN_OR_RAISE(auto min_value, ParquetMetrics::StatsValueToLiteral(
*column_desc, *iceberg_type, *stats,
/*is_min=*/true));
if (!lower_bound.has_value() || min_value < lower_bound.value()) {
lower_bound = std::move(min_value);
}

ICEBERG_ASSIGN_OR_RAISE(auto max_value,
StatsValueToLiteral(*column_desc, *iceberg_type, *stats,
ICEBERG_ASSIGN_OR_RAISE(auto max_value, ParquetMetrics::StatsValueToLiteral(
*column_desc, *iceberg_type, *stats,
/*is_min=*/false));
if (!upper_bound.has_value() || max_value > upper_bound.value()) {
upper_bound = std::move(max_value);
Expand Down Expand Up @@ -504,6 +463,47 @@ class CollectMetricsVisitor {

} // namespace

Result<Literal> ParquetMetrics::StatsValueToLiteral(
const ::parquet::ColumnDescriptor& column, const PrimitiveType& iceberg_type,
const ::parquet::Statistics& stats, bool is_min) {
switch (column.physical_type()) {
case ::parquet::Type::BOOLEAN:
return TypedStatsLiteral<::parquet::BoolStatistics>(
stats, is_min, [](bool value) { return Literal::Boolean(value); });
case ::parquet::Type::INT32:
return TypedStatsLiteral<::parquet::Int32Statistics>(
stats, is_min,
[&](int32_t value) { return Int32StatsLiteral(value, iceberg_type); });
case ::parquet::Type::INT64:
return TypedStatsLiteral<::parquet::Int64Statistics>(
stats, is_min,
[&](int64_t value) { return Int64StatsLiteral(value, iceberg_type); });
case ::parquet::Type::FLOAT:
return TypedStatsLiteral<::parquet::FloatStatistics>(
stats, is_min,
[&](float value) { return FloatStatsLiteral(value, iceberg_type); });
case ::parquet::Type::DOUBLE:
return TypedStatsLiteral<::parquet::DoubleStatistics>(
stats, is_min, [](double value) { return Literal::Double(value); });
case ::parquet::Type::BYTE_ARRAY:
return TypedStatsLiteral<::parquet::ByteArrayStatistics>(
stats, is_min, [&](const ::parquet::ByteArray& value) {
return BinaryStatsLiteral(BytesFromByteArray(value), iceberg_type);
});
case ::parquet::Type::FIXED_LEN_BYTE_ARRAY:
return TypedStatsLiteral<::parquet::FLBAStatistics>(
stats, is_min, [&](const ::parquet::FixedLenByteArray& value) {
return BinaryStatsLiteral(BytesFromFLBA(value, column.type_length()),
iceberg_type);
});
case ::parquet::Type::INT96:
case ::parquet::Type::UNDEFINED:
return NotSupported("Cannot convert Parquet statistics for physical type {}",
static_cast<int>(column.physical_type()));
}
std::unreachable();
}

Result<Metrics> ParquetMetrics::GetMetrics(
const Schema& schema, const ::parquet::SchemaDescriptor& parquet_schema,
const MetricsConfig& metrics_config, const ::parquet::FileMetaData& metadata,
Expand Down
6 changes: 6 additions & 0 deletions src/iceberg/parquet/parquet_metrics_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,12 @@ class ParquetMetrics {
public:
ParquetMetrics() = delete;

/// Convert one footer bound to an Iceberg literal, including supported promotions.
static Result<Literal> StatsValueToLiteral(const ::parquet::ColumnDescriptor& column,
const PrimitiveType& iceberg_type,
const ::parquet::Statistics& stats,
bool is_min);

/// \brief Compute file-level metrics from Parquet file metadata.
///
/// This function extracts metrics including row count, column sizes, value counts,
Expand Down
Loading
Loading