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
28 changes: 16 additions & 12 deletions cpp/src/io/parquet/experimental/page_index_filter.cu
Original file line number Diff line number Diff line change
Expand Up @@ -277,13 +277,13 @@ struct page_stats_caster : public stats_caster_base {

/**
* @brief Computes host side data including page row offsets, column chunk page offsets, and host
* columns containing page-level min, max and (optional) is_null statistics for a column
* columns containing page-level min, max and (optional) all-null statistics for a column
*
* @param schema_idx Column schema index
* @param dtype Column data type
* @param stream CUDA stream
* @return A tuple of page row offsets, column chunk page offsets, and host columns containing
* page-level min, max and (optional) is_null statistics
* page-level min, max and (optional) all-null statistics
*/
template <typename T>
[[nodiscard]] auto compute_host_data(cudf::size_type schema_idx,
Expand All @@ -300,11 +300,13 @@ struct page_stats_caster : public stats_caster_base {

auto const total_pages = col_chunk_page_offsets.back();

// Create host columns with page-level min, max and optionally is_null statistics
// Create host columns with page-level min, max and optionally all-null statistics. The
// all-null column is true only when every value in the page is null, false when none are, and
// null when only some are, which is what lets it answer both IS_NULL and IS NOT NULL.
host_column<T> min(total_pages, stream);
host_column<T> max(total_pages, stream);
std::optional<host_column<bool>> is_null;
if (has_is_null_operator) { is_null = host_column<bool>(total_pages, stream); }
std::optional<host_column<bool>> all_null;
if (has_is_null_operator) { all_null = host_column<bool>(total_pages, stream); }

// Compute timestamp scale factor for precision conversion
auto const ts_scale = [&] {
Expand Down Expand Up @@ -353,22 +355,24 @@ struct page_stats_caster : public stats_caster_base {
if (has_is_null_operator) {
// Check if the page is completely null
if (column_index.null_pages[page_idx]) {
is_null->val[column_page_idx] = true;
all_null->val[column_page_idx] = true;
return;
}
// Check if the page doesn't have a null count
if (not column_index.null_counts.has_value()) {
is_null->set_index(column_page_idx, std::nullopt, {});
all_null->set_index(column_page_idx, std::nullopt, {});
return;
}
// Use the null count to determine if the page is completely null
auto const page_row_count =
page_row_offsets[column_page_idx + 1] - page_row_offsets[column_page_idx];
auto const& null_count = column_index.null_counts.value()[page_idx];
if (null_count == page_row_count) {
is_null->val[column_page_idx] = false;
} else if (null_count > 0 and null_count < page_row_count) {
is_null->set_index(column_page_idx, std::nullopt, {});
if (null_count == 0) {
all_null->val[column_page_idx] = false;
} else if (null_count < page_row_count) {
all_null->set_index(column_page_idx, std::nullopt, {});
} else if (null_count == page_row_count) {
all_null->val[column_page_idx] = true;
} else {
CUDF_FAIL("Invalid null count");
}
Expand All @@ -381,7 +385,7 @@ struct page_stats_caster : public stats_caster_base {
std::move(col_chunk_page_offsets),
std::move(min),
std::move(max),
std::move(is_null)};
std::move(all_null)};
}

/**
Expand Down
104 changes: 104 additions & 0 deletions cpp/src/io/parquet/predicate_pushdown.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,110 @@

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 {
// Statistics without a null count say nothing about this chunk's nullability. The
// value array is allocated uninitialized and the null mask starts out all valid, so
// this entry has to be marked null; leaving it alone would let an uninitialized
// byte be read as an answer.
is_null->set_index(stats_idx, std::nullopt, {});
}
}
} 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
75 changes: 68 additions & 7 deletions cpp/src/io/parquet/stats_filter_helpers.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,29 @@

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

namespace {

/**
* @brief Maps a logical connective to its null-aware equivalent, returning any other operator as is
*
* A null in a statistics column says the writer did not record the statistic, never that the data
* is null, so the statistics expression is a three-valued predicate in which null means "unknown,
* keep this chunk". Three-valued logic is what propagates that: `false AND unknown` is false,
* because a chunk holding no row that can satisfy one conjunct cannot satisfy the conjunction
* whatever the other side turns out to be. The plain connectives instead return null whenever
* either side is null, which lets one absent statistic switch off pruning for the whole expression.
*/
[[nodiscard]] ast::ast_operator null_aware_operator(ast::ast_operator op)
{
switch (op) {
case ast::ast_operator::LOGICAL_AND: return ast::ast_operator::NULL_LOGICAL_AND;
case ast::ast_operator::LOGICAL_OR: return ast::ast_operator::NULL_LOGICAL_OR;
default: return op;
}
}

} // namespace

stats_columns_collector::stats_columns_collector(ast::expression const& expr,
cudf::size_type num_columns)
: _num_columns(num_columns)
Expand Down Expand Up @@ -76,6 +99,9 @@ std::reference_wrapper<ast::expression const> stats_columns_collector::visit(
op == ast_operator::LESS_EQUAL or op == ast_operator::GREATER or
op == ast_operator::GREATER_EQUAL) {
_columns_mask[col_ref->get_column_index()] = true;
// None of these can match a null, so their stats expressions consult the nullability column
// to rule out a chunk of nothing but nulls, which has no min or max to compare against.
_has_is_null_operator = true;
}
} else {
// Visit the operands and ignore any output as we only want to build the column mask
Expand All @@ -101,6 +127,28 @@ stats_expression_converter::stats_expression_converter(ast::expression const& ex
expr.accept(*this);
}

void stats_expression_converter::push_non_null_guard(size_type col_index,
ast::expression const& stats_expr)
{
using cudf::ast::ast_operator;

if (not std::cmp_equal(_stats_cols_per_column, 3)) { return; }

auto const& all_null =
_stats_expr.push(ast::column_reference{col_index * _stats_cols_per_column + 2});
// Answering "not entirely null" takes all three of the column's states, so a plain NOT will not
// do: its null state says the chunk holds both nulls and values, or that the writer recorded no
// null count, and both of those answer this question true. NOT alone answers it null and hands an
// unknown to a comparison that is in fact decisive.
auto const& not_all_null = _stats_expr.push(
ast::operation{ast_operator::NULL_LOGICAL_OR,
_stats_expr.push(ast::operation{ast_operator::IS_NULL, all_null}),
_stats_expr.push(ast::operation{ast_operator::NOT, all_null})});
// Null-aware so that the false this side pushes for an all-null chunk prunes it even though the
// min and max it lacks leave `stats_expr` unknown.
_stats_expr.push(ast::operation{ast_operator::NULL_LOGICAL_AND, not_all_null, stats_expr});
}

std::reference_wrapper<ast::expression const> stats_expression_converter::visit(
ast::operation const& expr)
{
Expand Down Expand Up @@ -203,35 +251,47 @@ std::reference_wrapper<ast::expression const> stats_expression_converter::visit(
_stats_expr.push(ast::column_reference{col_index * _stats_cols_per_column});
auto const& vmax =
_stats_expr.push(ast::column_reference{col_index * _stats_cols_per_column + 1});
_stats_expr.push(ast::operation{
ast::ast_operator::LOGICAL_AND,
// The two halves are separately optional in the statistics, so they are combined null-aware
// to keep whichever one is present decisive.
auto const& in_range = _stats_expr.push(ast::operation{
ast::ast_operator::NULL_LOGICAL_AND,
_stats_expr.push(ast::operation{ast_operator::GREATER_EQUAL, vmax, literal}),
_stats_expr.push(ast::operation{ast_operator::LESS_EQUAL, vmin, literal})});
// An all-null chunk has no min or max, so this range test is unknown there and would keep
// the chunk. The guard makes it prune instead.
push_non_null_guard(col_index, in_range);
break;
}
case ast_operator::NOT_EQUAL: {
auto const& vmin =
_stats_expr.push(ast::column_reference{col_index * _stats_cols_per_column});
auto const& vmax =
_stats_expr.push(ast::column_reference{col_index * _stats_cols_per_column + 1});
_stats_expr.push(
ast::operation{ast_operator::LOGICAL_OR,
// Null-aware for the same reason as the range test above: either half can be the one the
// statistics carry.
auto const& outside_range = _stats_expr.push(
ast::operation{ast_operator::NULL_LOGICAL_OR,
_stats_expr.push(ast::operation{ast_operator::NOT_EQUAL, vmin, vmax}),
_stats_expr.push(ast::operation{ast_operator::NOT_EQUAL, vmax, literal})});
// A null does not satisfy `!=` either, and an all-null chunk has no min or max to make this
// test decisive, so the guard prunes it.
push_non_null_guard(col_index, outside_range);
break;
}
case ast_operator::LESS: [[fallthrough]];
case ast_operator::LESS_EQUAL: {
auto const& vmin =
_stats_expr.push(ast::column_reference{col_index * _stats_cols_per_column});
_stats_expr.push(ast::operation{op, vmin, literal});
// An all-null chunk has no min, leaving this test unknown, so the guard prunes it.
push_non_null_guard(col_index, _stats_expr.push(ast::operation{op, vmin, literal}));
break;
}
case ast_operator::GREATER: [[fallthrough]];
case ast_operator::GREATER_EQUAL: {
auto const& vmax =
_stats_expr.push(ast::column_reference{col_index * _stats_cols_per_column + 1});
_stats_expr.push(ast::operation{op, vmax, literal});
// An all-null chunk has no max, leaving this test unknown, so the guard prunes it.
push_non_null_guard(col_index, _stats_expr.push(ast::operation{op, vmax, literal}));
break;
}
default: {
Expand All @@ -242,7 +302,8 @@ std::reference_wrapper<ast::expression const> stats_expression_converter::visit(
} // Visit operands and push expression for `expr op expr` form
else if (lhs_kind == operand_kind::EXPRESSION and rhs_kind == operand_kind::EXPRESSION) {
auto new_operands = visit_operands(expr.get_operands());
_stats_expr.push(ast::operation{op, new_operands.front(), new_operands.back()});
_stats_expr.push(
ast::operation{null_aware_operator(op), new_operands.front(), new_operands.back()});
} // Push _always_true for `col op col`, `expr op col`, `expr op lit` forms
else {
_stats_expr.push(ast::operation{ast_operator::IDENTITY, *_always_true});
Expand Down
26 changes: 24 additions & 2 deletions cpp/src/io/parquet/stats_filter_helpers.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -333,9 +333,13 @@ class stats_columns_collector : public ast::detail::expression_transformer {

/**
* @brief Return a boolean vector indicating input columns that can participate in stats based
* filtering
* filtering, and whether the stats table needs a per-column nullability column
*
* @return Boolean vector indicating input columns that can participate in stats based filtering
* The nullability column is needed by an `IS_NULL` operator, which is answered from it alone, and
* by any comparison against a literal, which uses it to rule out a chunk of nothing but nulls.
*
* @return Boolean vector indicating input columns that can participate in stats based filtering,
* and whether the nullability column is needed
*/
std::pair<thrust::host_vector<bool>, bool> get_stats_columns_mask() &&;

Expand Down Expand Up @@ -383,6 +387,24 @@ class stats_expression_converter : public stats_columns_collector {
thrust::host_vector<bool> get_stats_columns_mask() && = delete;

private:
/**
* @brief Push `not_all_null AND stats_expr` for a column, so that a chunk holding nothing but
* nulls fails a predicate that needs a non-null value to match
*
* A writer has no non-null value to compute min and max from for such a chunk, so it omits them
* and every min/max comparison evaluates to null, which keeps the chunk. The nullability
* statistic is decisive where min and max are absent: despite being built as `is_null`, it is
* true only when *every* value in the chunk is null, false when none are, and null when only some
* are or when the writer recorded no null count. Reading three states out of that column takes
* more than a `NOT`, since the null state answers "not entirely null" with a definite yes.
*
* Does nothing when the nullability column was not built, leaving `stats_expr` as the result.
*
* @param col_index Index of the column in the input table
* @param stats_expr Statistics expression to guard, already pushed onto the tree
*/
void push_non_null_guard(size_type col_index, ast::expression const& stats_expr);

ast::tree _stats_expr;
cudf::size_type _stats_cols_per_column;
std::unique_ptr<cudf::numeric_scalar<bool>> _always_true_scalar;
Expand Down
Loading
Loading