Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
31 commits
Select commit Hold shift + click to select a range
ec40676
expose min/max column-chunk statistics in pylibcudf
rjzamora Aug 14, 2026
6c60795
use min/max_encoded
rjzamora Aug 14, 2026
e71e4df
add column_chunk_bounds helper
rjzamora Aug 14, 2026
a039508
Merge remote-tracking branch 'upstream/main' into parquet-minmax-stat…
rjzamora Aug 14, 2026
9de5438
hopefully fix CI
rjzamora Aug 14, 2026
6ea7be5
fix linting
rjzamora Aug 14, 2026
8c871cb
Merge remote-tracking branch 'upstream/main' into parquet-minmax-stat…
rjzamora Aug 14, 2026
25f4e45
try fixing lint error
rjzamora Aug 14, 2026
996ab71
Merge remote-tracking branch 'upstream/main' into parquet-minmax-stat…
rjzamora Aug 14, 2026
466cb54
fix clang
rjzamora Aug 17, 2026
46fcaae
Merge remote-tracking branch 'upstream/main' into parquet-minmax-stat…
rjzamora Aug 17, 2026
df98442
fix build
rjzamora Aug 17, 2026
200792c
fix comment
rjzamora Aug 17, 2026
c61b4cb
trigger CI
rjzamora Aug 17, 2026
f688c0e
fix
rjzamora Aug 17, 2026
e0b16ab
Merge remote-tracking branch 'upstream/main' into parquet-minmax-stat…
rjzamora Aug 17, 2026
f920ee1
align with main
rjzamora Aug 17, 2026
00bee2d
fix fallback for non-signed
rjzamora Aug 18, 2026
0d4141f
Merge remote-tracking branch 'upstream/main' into parquet-minmax-stat…
rjzamora Aug 18, 2026
4adf9d7
Merge branch 'main' into parquet-minmax-stats-helper
rjzamora Aug 18, 2026
6a3b6c6
Merge branch 'main' into parquet-minmax-stats-helper
rjzamora Aug 18, 2026
a936e56
adress smaller code-review comments
rjzamora Aug 19, 2026
1e62b62
Merge remote-tracking branch 'upstream/main' into parquet-minmax-stat…
rjzamora Aug 19, 2026
3c8fab7
Merge remote-tracking branch 'upstream/main' into parquet-minmax-stat…
rjzamora Aug 20, 2026
82c6e71
return single Table
rjzamora Aug 20, 2026
8c58992
Merge remote-tracking branch 'upstream/main' into parquet-minmax-stat…
rjzamora Aug 20, 2026
a6d0cdb
fix cython
rjzamora Aug 20, 2026
418e256
take haseeb's suggestion
rjzamora Aug 20, 2026
8d48621
Merge remote-tracking branch 'upstream/main' into parquet-minmax-stat…
rjzamora Aug 20, 2026
baf8c29
Merge branch 'main' into parquet-minmax-stats-helper
rjzamora Aug 21, 2026
8316cdd
address review
rjzamora Aug 21, 2026
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
35 changes: 35 additions & 0 deletions cpp/include/cudf/io/parquet_metadata.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,14 @@
#include <cudf/io/datasource.hpp>
#include <cudf/io/parquet_schema.hpp>
#include <cudf/io/types.hpp>
#include <cudf/table/table.hpp>
#include <cudf/utilities/default_stream.hpp>
#include <cudf/utilities/export.hpp>
#include <cudf/utilities/memory_resource.hpp>

#include <memory>
#include <span>
#include <string>
#include <string_view>
#include <vector>

Expand Down Expand Up @@ -294,6 +299,36 @@ parquet_metadata read_parquet_metadata(source_info const& src_info);
std::vector<parquet::FileMetaData> read_parquet_footers(
std::span<std::unique_ptr<cudf::io::datasource> const> sources);

/**
* @brief Decode parquet column-chunk min/max statistics for selected leaf columns.
*
* Missing min/max statistics are represented as nulls in the corresponding output column. Parquet
* min/max exactness flags are not interpreted by this function. The requested column names are
* resolved against each file's schema. The returned table contains one row per source row group.
* Column 0 is the source file index, column 1 is the file-local row-group index, and subsequent
* columns are min/max pairs in the order of ``column_names``.
*
* @ingroup io_readers
*
* @param parquet_metadatas Parquet file metadata, one per source
* @param column_names Dotted leaf-column paths to decode statistics for
* @param stream CUDA stream used for device memory operations
* @param mr Memory resources to use for device memory allocation
* @return Table of row-group identifiers and decoded min/max bounds. For requested column
* ``column_names[i]``, the min column is at ``2 + 2 * i`` and the max column is at
* ``3 + 2 * i``.
*
* @throw std::invalid_argument If a requested leaf-column path is missing or ambiguous.
* @throw std::invalid_argument If a requested column has unsupported or compound statistics dtype.
* @throw std::invalid_argument If a requested column has mismatching statistics dtype across
* sources.
*/
std::unique_ptr<table> read_parquet_column_chunk_bounds(
std::span<parquet::FileMetaData const> parquet_metadatas,
std::span<std::string const> column_names,
cuda::stream_ref stream = cudf::get_default_stream(),
cudf::memory_resources mr = cudf::get_current_device_resource_ref());

/** @} */ // end of group
} // namespace io
} // namespace CUDF_EXPORT cudf
17 changes: 17 additions & 0 deletions cpp/src/io/functions.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -683,6 +683,23 @@ std::vector<parquet::FileMetaData> read_parquet_footers(
return detail_parquet::read_parquet_footers(sources);
}

std::unique_ptr<table> read_parquet_column_chunk_bounds(
std::span<parquet::FileMetaData const> parquet_metadatas,
std::span<std::string const> column_names,
cuda::stream_ref stream,
cudf::memory_resources mr)
{
CUDF_FUNC_RANGE();

auto metadata = detail_parquet::aggregate_reader_metadata{
std::vector<parquet::FileMetaData>{parquet_metadatas.begin(), parquet_metadatas.end()},
false, // use_arrow_schema
false // has_cols_from_mismatched_srcs
};

return metadata.read_column_chunk_bounds(column_names, stream, mr.get_output_mr());
}

/**
* @copydoc cudf::io::merge_row_group_metadata
*/
Expand Down
110 changes: 5 additions & 105 deletions cpp/src/io/parquet/predicate_pushdown.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -5,17 +5,16 @@

#include "expression_transform_helpers.hpp"
#include "reader_impl_helpers.hpp"
#include "row_group_stats_helpers.hpp"
#include "stats_filter_helpers.hpp"
#include "timestamp_utils.cuh"

#include <cudf/column/column_factories.hpp>
#include <cudf/detail/iterator.cuh>
#include <cudf/detail/transform.hpp>
#include <cudf/detail/utilities/vector_factories.hpp>
#include <cudf/table/table.hpp>
#include <cudf/utilities/error.hpp>
#include <cudf/utilities/memory_resource.hpp>
#include <cudf/utilities/span.hpp>
#include <cudf/utilities/traits.hpp>
#include <cudf/utilities/type_dispatcher.hpp>

#include <cuda/iterator>
Expand All @@ -28,104 +27,6 @@

namespace cudf::io::parquet::detail {

namespace {

/**
* @brief Converts column chunk statistics to 2 device columns - min, max values.
*
* Each column's number of rows equals the total number of row groups.
*
*/
struct row_group_stats_caster : public stats_caster_base {
size_type total_row_groups;
std::vector<metadata> const& per_file_metadata;
host_span<std::vector<size_type> const> row_group_indices;
bool has_is_null_operator;

// Creates device columns from column statistics (min, max)
template <typename T>
std::
tuple<std::unique_ptr<column>, std::unique_ptr<column>, std::optional<std::unique_ptr<column>>>
operator()(host_span<int const> per_source_schema_indices,
cudf::data_type dtype,
cuda::stream_ref stream,
rmm::device_async_resource_ref mr) const
{
// List, Struct, Dictionary types are not supported
if constexpr (cudf::is_compound<T>() && !std::is_same_v<T, string_view>) {
CUDF_FAIL("Compound types do not have statistics");
} else {
host_column<T> min(total_row_groups, stream);
host_column<T> max(total_row_groups, stream);
std::optional<host_column<bool>> is_null;
if (has_is_null_operator) { is_null = host_column<bool>(total_row_groups, stream); }

size_type stats_idx = 0;
for (size_t src_idx = 0; src_idx < row_group_indices.size(); ++src_idx) {
auto const mapped_schema_idx = per_source_schema_indices[src_idx];
// Compute timestamp scale factor for precision conversion from the mapped source schema.
auto const ts_scale = [&] {
if constexpr (cudf::is_timestamp<T>()) {
auto const& schema = per_file_metadata[src_idx].schema[mapped_schema_idx];
return calc_timestamp_scale(schema.logical_type, static_cast<int32_t>(T::period::den));
}
return 0;
}();

for (auto const rg_idx : row_group_indices[src_idx]) {
auto const& row_group = per_file_metadata[src_idx].row_groups[rg_idx];
auto col = std::find_if(row_group.columns.begin(),
row_group.columns.end(),
[mapped_schema_idx](ColumnChunk const& col) {
return col.schema_idx == mapped_schema_idx;
});
if (col != std::end(row_group.columns)) {
auto const& colchunk = *col;
// To support deprecated min, max fields.
auto const& min_value = colchunk.meta_data.statistics.min_value.has_value()
? colchunk.meta_data.statistics.min_value
: colchunk.meta_data.statistics.min;
auto const& max_value = colchunk.meta_data.statistics.max_value.has_value()
? colchunk.meta_data.statistics.max_value
: colchunk.meta_data.statistics.max;
// translate binary data to Type then to <T>
min.set_index(stats_idx, min_value, colchunk.meta_data.type, ts_scale);
max.set_index(stats_idx, max_value, colchunk.meta_data.type, ts_scale);
// Check the nullability of this column chunk
if (has_is_null_operator) {
if (colchunk.meta_data.statistics.null_count.has_value()) {
auto const& null_count = colchunk.meta_data.statistics.null_count.value();
if (null_count == 0) {
is_null->val[stats_idx] = false;
} else if (null_count < colchunk.meta_data.num_values) {
is_null->set_index(stats_idx, std::nullopt, {});
} else if (null_count == colchunk.meta_data.num_values) {
is_null->val[stats_idx] = true;
} else {
CUDF_FAIL("Invalid null count");
}
}
}
} else {
// Marking it null, if column present in row group
min.set_index(stats_idx, std::nullopt, {});
max.set_index(stats_idx, std::nullopt, {});
if (has_is_null_operator) { is_null->set_index(stats_idx, std::nullopt, {}); }
}
stats_idx++;
}
};
return {min.to_device(dtype, stream, mr),
max.to_device(dtype, stream, mr),
has_is_null_operator ? std::make_optional(is_null->to_device(
data_type{cudf::type_id::BOOL8}, stream, mr))
: std::nullopt};
}
}
};

} // namespace

bool aggregate_reader_metadata::any_row_group_stats_available(
host_span<std::vector<size_type> const> input_row_group_indices,
host_span<int const> filter_column_schemas) const
Expand Down Expand Up @@ -286,10 +187,9 @@ aggregate_reader_metadata::filter_row_groups(
: total_row_groups;

// Span of row groups to apply bloom filtering on.
auto const bloom_filter_input_row_groups =
stats_filtered_row_groups.has_value()
? host_span<std::vector<size_type> const>(stats_filtered_row_groups.value())
: input_row_group_indices;
auto const bloom_filter_input_row_groups = stats_filtered_row_groups.has_value()
? stats_filtered_row_groups.value()
: input_row_group_indices;

// Collect equality literals for each input table column for bloom filtering
auto const equality_literals =
Expand Down
114 changes: 113 additions & 1 deletion cpp/src/io/parquet/reader_impl_helpers.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -12,13 +12,19 @@
#include "ipc/Message_generated.h"
#include "ipc/Schema_generated.h"
#include "parquet_common.hpp"
#include "row_group_stats_helpers.hpp"

#include <cudf/column/column.hpp>
#include <cudf/detail/nvtx/ranges.hpp>
#include <cudf/detail/utilities/host_memory.hpp>
#include <cudf/detail/utilities/host_worker_pool.hpp>
#include <cudf/io/parquet_io_utils.hpp>
#include <cudf/io/parquet_metadata.hpp>
#include <cudf/io/parquet_schema.hpp>
#include <cudf/logger.hpp>
#include <cudf/table/table.hpp>
#include <cudf/utilities/error.hpp>
#include <cudf/utilities/traits.hpp>
#include <cudf/utilities/type_dispatcher.hpp>

#include <cuda/iterator>
#include <cuda/numeric>
Expand All @@ -29,9 +35,11 @@
#include <format>
#include <functional>
#include <future>
#include <iterator>
#include <numeric>
#include <optional>
#include <regex>
#include <span>
#include <string_view>
#include <utility>

Expand Down Expand Up @@ -79,6 +87,43 @@ namespace flatbuf = cudf::io::parquet::flatbuf;

namespace {

[[nodiscard]] int find_leaf_schema_index(std::span<SchemaElement const> schema_tree,
std::string_view column_name)
{
auto found = std::optional<int>{};
for (auto idx = 1; std::cmp_less(idx, schema_tree.size()); ++idx) {
if (not schema_tree[idx].children_idx.empty()) { continue; }
if (column_path_from_index(schema_tree, idx) != column_name) { continue; }
// Parquet stores names per schema node and does not enforce unique dotted leaf paths.
// Reject ambiguity here instead of returning the first matching schema element.
CUDF_EXPECTS(not found.has_value(),
Comment thread
rjzamora marked this conversation as resolved.
std::string{"Ambiguous parquet leaf column path: "} + std::string{column_name},
std::invalid_argument);
found = idx;
}
CUDF_EXPECTS(found.has_value(),
std::string{"Parquet leaf column path not found: "} + std::string{column_name},
std::invalid_argument);
return found.value();
}

[[nodiscard]] data_type statistics_dtype(SchemaElement const& schema)
{
auto const dtype = to_data_type(to_type_id(schema,
false, // strings_to_categorical
type_id::EMPTY, // timestamp_type_id
type_id::EMPTY),
schema);
CUDF_EXPECTS(dtype.id() != type_id::EMPTY,
std::string{"Unsupported parquet statistics dtype for column: "} + schema.name,
std::invalid_argument);
CUDF_EXPECTS(
not cudf::is_compound(dtype) or dtype.id() == type_id::STRING,
std::string{"Compound parquet statistics are not supported for column: "} + schema.name,
std::invalid_argument);
return dtype;
}

/**
* @brief Computes the total number of row groups in input span of row group indices
*/
Expand Down Expand Up @@ -1366,6 +1411,73 @@ aggregate_reader_metadata::get_column_chunk_metadata() const
return column_chunk_metadata;
}

std::unique_ptr<table> aggregate_reader_metadata::read_column_chunk_bounds(
std::span<std::string const> column_names,
cuda::stream_ref stream,
rmm::device_async_resource_ref mr) const
{
CUDF_EXPECTS(column_names.empty() or not per_file_metadata.empty(),
"Cannot decode parquet column-chunk bounds without source metadata",
std::invalid_argument);

auto const total_row_groups = get_num_row_groups();
auto const num_row_groups_per_file = get_num_row_groups_per_file();
auto num_row_groups_per_source = std::vector<std::size_t>{};
num_row_groups_per_source.reserve(num_row_groups_per_file.size());
std::transform(num_row_groups_per_file.begin(),
num_row_groups_per_file.end(),
std::back_inserter(num_row_groups_per_source),
[](auto count) { return static_cast<std::size_t>(count); });

auto input_row_group_indices = std::vector<std::vector<size_type>>(per_file_metadata.size());
for (auto src_idx = size_type{0}; std::cmp_less(src_idx, per_file_metadata.size()); ++src_idx) {
Comment thread
mhaseeb123 marked this conversation as resolved.
auto& source_row_group_indices = input_row_group_indices[src_idx];
source_row_group_indices.resize(num_row_groups_per_file[src_idx]);
std::iota(source_row_group_indices.begin(), source_row_group_indices.end(), size_type{0});
}

std::vector<std::unique_ptr<column>> columns;
columns.reserve(2 + 2 * column_names.size());
auto file_indices = synthesize_source_index_column(num_row_groups_per_source, stream, mr);
auto row_group_indices = synthesize_row_group_index_column(file_indices->view(), stream, mr);
columns.push_back(std::move(file_indices));
columns.push_back(std::move(row_group_indices));

row_group_stats_caster const stats_col{.total_row_groups = total_row_groups,
.per_file_metadata = per_file_metadata,
.row_group_indices = input_row_group_indices,
.has_is_null_operator = false};

for (auto const& column_name : column_names) {
auto per_source_schema_indices = std::vector<int>(per_file_metadata.size());
auto dtype = data_type{type_id::EMPTY};

for (auto src_idx = size_type{0}; std::cmp_less(src_idx, per_file_metadata.size()); ++src_idx) {
auto const& schema_tree = get_schema_tree(src_idx);
auto const schema_idx = find_leaf_schema_index(schema_tree, column_name);
auto const source_dtype = statistics_dtype(schema_tree[schema_idx]);

if (src_idx == 0) {
dtype = source_dtype;
} else {
CUDF_EXPECTS(
source_dtype == dtype,
std::string{"Mismatching parquet statistics dtype across sources for column: "} +
column_name,
std::invalid_argument);
}
per_source_schema_indices[src_idx] = schema_idx;
}

auto [min_col, max_col, _] = cudf::type_dispatcher<dispatch_storage_type>(
dtype, stats_col, per_source_schema_indices, dtype, stream, mr);
columns.push_back(std::move(min_col));
columns.push_back(std::move(max_col));
}

return std::make_unique<table>(std::move(columns));
}

bool aggregate_reader_metadata::is_schema_index_mapped(int schema_idx, int src_idx) const
{
// Check if schema_idx or src_idx is invalid
Expand Down
Loading
Loading