diff --git a/cpp/src/io/parquet/experimental/page_index_filter.cu b/cpp/src/io/parquet/experimental/page_index_filter.cu index 7ec4aa859f0..a4596553b65 100644 --- a/cpp/src/io/parquet/experimental/page_index_filter.cu +++ b/cpp/src/io/parquet/experimental/page_index_filter.cu @@ -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 [[nodiscard]] auto compute_host_data(cudf::size_type schema_idx, @@ -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 min(total_pages, stream); host_column max(total_pages, stream); - std::optional> is_null; - if (has_is_null_operator) { is_null = host_column(total_pages, stream); } + std::optional> all_null; + if (has_is_null_operator) { all_null = host_column(total_pages, stream); } // Compute timestamp scale factor for precision conversion auto const ts_scale = [&] { @@ -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"); } @@ -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)}; } /** diff --git a/cpp/src/io/parquet/predicate_pushdown.cpp b/cpp/src/io/parquet/predicate_pushdown.cpp index 44fa83badc7..378ce7f216c 100644 --- a/cpp/src/io/parquet/predicate_pushdown.cpp +++ b/cpp/src/io/parquet/predicate_pushdown.cpp @@ -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 const& per_file_metadata; + host_span const> row_group_indices; + bool has_is_null_operator; + + // Creates device columns from column statistics (min, max) + template + std:: + tuple, std::unique_ptr, std::optional>> + operator()(host_span 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() && !std::is_same_v) { + CUDF_FAIL("Compound types do not have statistics"); + } else { + host_column min(total_row_groups, stream); + host_column max(total_row_groups, stream); + std::optional> is_null; + if (has_is_null_operator) { is_null = host_column(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()) { + auto const& schema = per_file_metadata[src_idx].schema[mapped_schema_idx]; + return calc_timestamp_scale(schema.logical_type, static_cast(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 + 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 const> input_row_group_indices, host_span filter_column_schemas) const diff --git a/cpp/src/io/parquet/stats_filter_helpers.cpp b/cpp/src/io/parquet/stats_filter_helpers.cpp index fb5dbd3e1e1..53bfe23de8f 100644 --- a/cpp/src/io/parquet/stats_filter_helpers.cpp +++ b/cpp/src/io/parquet/stats_filter_helpers.cpp @@ -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) @@ -76,6 +99,9 @@ std::reference_wrapper 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 @@ -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 stats_expression_converter::visit( ast::operation const& expr) { @@ -203,10 +251,15 @@ std::reference_wrapper 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: { @@ -214,24 +267,31 @@ std::reference_wrapper 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_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: { @@ -242,7 +302,8 @@ std::reference_wrapper 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}); diff --git a/cpp/src/io/parquet/stats_filter_helpers.hpp b/cpp/src/io/parquet/stats_filter_helpers.hpp index ed7ac756bd2..9b6d4579030 100644 --- a/cpp/src/io/parquet/stats_filter_helpers.hpp +++ b/cpp/src/io/parquet/stats_filter_helpers.hpp @@ -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, bool> get_stats_columns_mask() &&; @@ -383,6 +387,24 @@ class stats_expression_converter : public stats_columns_collector { thrust::host_vector 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> _always_true_scalar; diff --git a/cpp/tests/io/parquet_reader_test.cpp b/cpp/tests/io/parquet_reader_test.cpp index 6e55ccbd5ca..98b53688a6d 100644 --- a/cpp/tests/io/parquet_reader_test.cpp +++ b/cpp/tests/io/parquet_reader_test.cpp @@ -2898,6 +2898,113 @@ TEST_F(ParquetReaderTest, FilterNoStats) CUDF_TEST_EXPECT_TABLES_EQUAL(expected->view(), result); } +// Filter on a column whose row groups differ in nullability, which is what makes a statistic +// indecisive: a chunk of nothing but nulls has no min or max at all, and a chunk holding both nulls +// and values has a null count that says neither of those things. +TEST_F(ParquetReaderTest, FilterNullableStats) +{ + auto constexpr num_input_row_groups = 3; + + auto const filepath = temp_env->get_temp_filepath("FilterNullableStats.parquet"); + + // Three row groups of three rows. Column `a` holds some nulls in the first, nothing but nulls in + // the second and none in the third, so its nullability statistic takes each of its three states. + // Column `b` is never null and holds values far below the literal compared against it below. + { + auto const a0 = + cudf::test::fixed_width_column_wrapper({10, 0, 20}, {true, false, true}); + auto const a1 = + cudf::test::fixed_width_column_wrapper({0, 0, 0}, {false, false, false}); + auto const a2 = cudf::test::fixed_width_column_wrapper({100, 200, 300}); + auto const b0 = cudf::test::fixed_width_column_wrapper({1, 1, 1}); + auto const b1 = cudf::test::fixed_width_column_wrapper({2, 2, 2}); + auto const b2 = cudf::test::fixed_width_column_wrapper({3, 3, 3}); + auto const t0 = cudf::table_view{{a0, b0}}; + auto const t1 = cudf::table_view{{a1, b1}}; + auto const t2 = cudf::table_view{{a2, b2}}; + + auto const options = + cudf::io::chunked_parquet_writer_options::builder(cudf::io::sink_info{filepath}) + .metadata(cudf::io::table_input_metadata(t0)) + .build(); + + cudf::io::chunked_parquet_writer writer(options); + writer.write(t0); + writer.write(t1); + writer.write(t2); + writer.close(); + } + + auto const test_predicate_pushdown = [&](cudf::ast::operation const& filter, + cudf::size_type expected_filtered_row_groups, + cudf::size_type expected_num_rows) { + auto const options = cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepath}) + .filter(filter) + .build(); + + auto const result = cudf::io::read_parquet(options); + + EXPECT_EQ(result.metadata.num_input_row_groups, num_input_row_groups); + EXPECT_TRUE(result.metadata.num_row_groups_after_stats_filter.has_value()); + EXPECT_EQ(result.metadata.num_row_groups_after_stats_filter.value(), + expected_filtered_row_groups); + EXPECT_EQ(result.tbl->num_rows(), expected_num_rows); + }; + + auto const a_ref = cudf::ast::column_reference(0); + auto const b_ref = cudf::ast::column_reference(1); + + auto scalar_10 = cudf::numeric_scalar(10, true); + auto scalar_20 = cudf::numeric_scalar(20, true); + auto scalar_50 = cudf::numeric_scalar(50, true); + auto scalar_60 = cudf::numeric_scalar(60, true); + auto scalar_5 = cudf::numeric_scalar(5, true); + auto const literal_10 = cudf::ast::literal(scalar_10); + auto const literal_20 = cudf::ast::literal(scalar_20); + auto const literal_50 = cudf::ast::literal(scalar_50); + auto const literal_60 = cudf::ast::literal(scalar_60); + auto const literal_5 = cudf::ast::literal(scalar_5); + + auto const a_is_null = cudf::ast::operation(cudf::ast::ast_operator::IS_NULL, a_ref); + auto const a_ge_10 = + cudf::ast::operation(cudf::ast::ast_operator::GREATER_EQUAL, a_ref, literal_10); + auto const a_le_20 = cudf::ast::operation(cudf::ast::ast_operator::LESS_EQUAL, a_ref, literal_20); + auto const a_ge_50 = + cudf::ast::operation(cudf::ast::ast_operator::GREATER_EQUAL, a_ref, literal_50); + auto const a_le_60 = cudf::ast::operation(cudf::ast::ast_operator::LESS_EQUAL, a_ref, literal_60); + auto const b_gt_5 = cudf::ast::operation(cudf::ast::ast_operator::GREATER, b_ref, literal_5); + + // Filter: IS_NULL(a). The all-null row group answers this yes and the partly null one cannot + // answer it at all, so both are kept and only the row group with no nulls is ruled out. + test_predicate_pushdown(a_is_null, 2, 4); + + // Filter: a >= 10 AND a <= 20. RG 0 passes on its values, which the nulls it also holds must not + // count against; RG 1 holds nothing a comparison can match; RG 2's min of 100 rules it out. + { + auto const filter = + cudf::ast::operation(cudf::ast::ast_operator::LOGICAL_AND, a_ge_10, a_le_20); + test_predicate_pushdown(filter, 1, 2); + } + + // Filter: a >= 50 AND a <= 60 — matches no row group. RG 0's max of 20 rules it out even though + // its other conjunct is indecisive there, which is the case a conjunction that is not null-aware + // gets wrong: it would carry the indecisive side up and keep the row group. + { + auto const filter = + cudf::ast::operation(cudf::ast::ast_operator::LOGICAL_AND, a_ge_50, a_le_60); + test_predicate_pushdown(filter, 0, 0); + } + + // Filter: IS_NULL(a) AND b > 5 — matches no row group, since `b` reaches only 3. The `IS_NULL` + // side is indecisive on RG 0 and decides nothing on its own anywhere, so this pins that one + // decisive conjunct is enough to prune whatever the other side says. + { + auto const filter = + cudf::ast::operation(cudf::ast::ast_operator::LOGICAL_AND, a_is_null, b_gt_5); + test_predicate_pushdown(filter, 0, 0); + } +} + // Filter for float column with NaN values TEST_F(ParquetReaderTest, FilterFloatNAN) { @@ -4399,12 +4506,14 @@ void filter_unary_operation_typed_test() auto const ref_not_expr1 = cudf::ast::operation(cudf::ast::ast_operator::NOT, ref_expr1); auto const ref_expr2 = cudf::ast::operation(cudf::ast::ast_operator::IS_NULL, col_ref_0); - // col0 < 100 AND IS_NULL(col0) + // col0 < 100 AND IS_NULL(col0). No row satisfies this, since a null is not less than anything, + // so every row group is ruled out: the all-null one by the comparison, which needs a non-null + // value to match, and the rest by `IS_NULL` against statistics that count no nulls. auto filter_expression = cudf::ast::operation(cudf::ast::ast_operator::LOGICAL_AND, expr1, expr2); auto ref_filter = cudf::ast::operation(cudf::ast::ast_operator::LOGICAL_AND, ref_expr1, ref_expr2); - auto constexpr expected_filtered_row_groups_with_unary_and = 1; + auto constexpr expected_filtered_row_groups_with_unary_and = 0; test_predicate_pushdown(filter_expression, ref_filter, expected_total_row_groups,