From 33f2306cd0887fbc3d3e536d45cb2bddb210d148 Mon Sep 17 00:00:00 2001 From: jsonwu Date: Thu, 6 Nov 2025 13:27:32 +0800 Subject: [PATCH 1/4] fix a bug of materialize boolean --- cpp/src/arrow/acero/sorted_merge_node_test.cc | 110 ++++++++++++++++++ .../acero/unmaterialized_table_internal.h | 7 +- 2 files changed, 116 insertions(+), 1 deletion(-) diff --git a/cpp/src/arrow/acero/sorted_merge_node_test.cc b/cpp/src/arrow/acero/sorted_merge_node_test.cc index 82b630420c4..8832ddd4a78 100644 --- a/cpp/src/arrow/acero/sorted_merge_node_test.cc +++ b/cpp/src/arrow/acero/sorted_merge_node_test.cc @@ -16,16 +16,23 @@ // under the License. #include +#include +#include #include "arrow/acero/exec_plan.h" #include "arrow/acero/map_node.h" #include "arrow/acero/options.h" #include "arrow/acero/test_nodes.h" #include "arrow/array/builder_base.h" +#include "arrow/array/builder_primitive.h" #include "arrow/array/concatenate.h" #include "arrow/compute/ordering.h" +#include "arrow/dataset/dataset.h" +#include "arrow/dataset/scanner.h" +#include "arrow/record_batch.h" #include "arrow/result.h" #include "arrow/scalar.h" +#include "arrow/status.h" #include "arrow/table.h" #include "arrow/testing/generator.h" #include "arrow/testing/gtest_util.h" @@ -83,4 +90,107 @@ TEST(SortedMergeNode, Basic) { AssertArraysEqual(*expected_ts, *output_ts); } +TEST(SortedMergeNode, TestSortedMergeTwoInputsWithBool) { + const int64_t row_count = (16 << 10); // 16k rows per input + + // Create schema with int column A and bool column B + auto test_schema = arrow::schema( + {arrow::field("col_a", arrow::int32()), arrow::field("col_b", arrow::boolean())}); + + // Helper lambda to create table with specific pattern + auto create_test_scanner = [&](int64_t cnt, int offset) -> arrow::Result { + // Create column A (int) - values from offset to offset+cnt-1 + arrow::Int32Builder col_a_builder; + std::vector col_a_values; + col_a_values.reserve(cnt); + for (int64_t i = 0; i < cnt; ++i) { + col_a_values.push_back(static_cast(offset + i)); + } + ARROW_RETURN_NOT_OK(col_a_builder.AppendValues(col_a_values)); + std::shared_ptr col_a_arr; + ARROW_RETURN_NOT_OK(col_a_builder.Finish(&col_a_arr)); + + // Create column B (bool) - pattern: true if col_a % 5 == 0, false otherwise + arrow::BooleanBuilder col_b_builder; + for (int64_t i = 0; i < cnt; ++i) { + int32_t a_value = offset + i; + bool b_value = (a_value % 5 == 0); + ARROW_RETURN_NOT_OK(col_b_builder.Append(b_value)); + } + std::shared_ptr col_b_arr; + ARROW_RETURN_NOT_OK(col_b_builder.Finish(&col_b_arr)); + + auto table = arrow::Table::Make(test_schema, {col_a_arr, col_b_arr}); + auto table_source = + Declaration("table_source", TableSourceNodeOptions(table, row_count / 16)); + return table_source; + }; + + ASSERT_OK_AND_ASSIGN(auto source1, create_test_scanner(row_count, 0)); + ASSERT_OK_AND_ASSIGN(auto source2, create_test_scanner(row_count, 8192)); + + // Create sorted merge by column A + auto ops = OrderByNodeOptions(compute::Ordering({compute::SortKey("col_a")})); + Declaration sorted_merge{"sorted_merge", {source1, source2}, ops}; + + // Execute plan and collect result + ASSERT_OK_AND_ASSIGN(auto result_table, + arrow::acero::DeclarationToTable(sorted_merge, false)); + + ASSERT_TRUE(result_table != nullptr); + + // Verify results + auto col_a = result_table->GetColumnByName("col_a"); + auto col_b = result_table->GetColumnByName("col_b"); + ASSERT_TRUE(col_a != nullptr); + ASSERT_TRUE(col_b != nullptr); + + // Verify sorting and bool values + int32_t last_a_value = std::numeric_limits::min(); + int64_t total_rows_checked = 0; + int64_t true_count = 0; + int64_t false_count = 0; + + for (int i = 0; i < col_a->num_chunks(); i++) { + auto a_chunk = std::static_pointer_cast(col_a->chunk(i)); + auto b_chunk = std::static_pointer_cast(col_b->chunk(i)); + + ASSERT_EQ(a_chunk->length(), b_chunk->length()) + << "Column A and B must have same length in chunk " << i; + + for (int64_t j = 0; j < a_chunk->length(); j++) { + ASSERT_FALSE(a_chunk->IsNull(j)) << "Column A should not have null values"; + ASSERT_FALSE(b_chunk->IsNull(j)) << "Column B should not have null values"; + + int32_t a_value = a_chunk->Value(j); + bool b_value = b_chunk->Value(j); + + // Verify sorting by column A + ASSERT_GE(a_value, last_a_value) + << "Values not sorted at chunk " << i << ", row " << j + << ": current=" << a_value << ", last=" << last_a_value; + last_a_value = a_value; + + // Verify bool value correctness: should be true if a_value % 3 == 0 + bool expected_b_value = (a_value % 5 == 0); + ASSERT_EQ(b_value, expected_b_value) + << "Bool value incorrect at chunk " << i << ", row " << j + << ": col_a=" << a_value << ", col_b=" << b_value + << ", expected=" << expected_b_value; + + if (b_value) { + true_count++; + } else { + false_count++; + } + total_rows_checked++; + } + } + + ASSERT_EQ(last_a_value, 24575); + + ASSERT_EQ(total_rows_checked, row_count * 2) + << "Expected " << row_count << " unique rows after merge"; +} + } // namespace arrow::acero diff --git a/cpp/src/arrow/acero/unmaterialized_table_internal.h b/cpp/src/arrow/acero/unmaterialized_table_internal.h index 86b1a763a60..2c6b055fc43 100644 --- a/cpp/src/arrow/acero/unmaterialized_table_internal.h +++ b/cpp/src/arrow/acero/unmaterialized_table_internal.h @@ -172,7 +172,12 @@ class UnmaterializedCompositeTable { builder.UnsafeAppendNull(); return Status::OK(); } - builder.UnsafeAppend(bit_util::GetBit(source->template GetValues(1), row)); + + int64_t array_offset = (source->offset + row) / 8; + int64_t bit_offset = (source->offset + row) % 8; + + builder.UnsafeAppend(arrow::bit_util::GetBit( + source->template GetValues(1, array_offset), bit_offset)); return Status::OK(); } From e1731e059b4d4252ff14a093e049d9f576a96d5f Mon Sep 17 00:00:00 2001 From: Nic Crane Date: Tue, 29 Sep 2026 14:10:13 +0100 Subject: [PATCH 2/4] Fix -Werror failures in sorted merge bool test; cover unaligned offsets and nulls --- cpp/src/arrow/acero/sorted_merge_node_test.cc | 56 ++++++++++--------- 1 file changed, 31 insertions(+), 25 deletions(-) diff --git a/cpp/src/arrow/acero/sorted_merge_node_test.cc b/cpp/src/arrow/acero/sorted_merge_node_test.cc index 8832ddd4a78..b3f26356418 100644 --- a/cpp/src/arrow/acero/sorted_merge_node_test.cc +++ b/cpp/src/arrow/acero/sorted_merge_node_test.cc @@ -16,6 +16,7 @@ // under the License. #include +#include #include #include @@ -27,8 +28,6 @@ #include "arrow/array/builder_primitive.h" #include "arrow/array/concatenate.h" #include "arrow/compute/ordering.h" -#include "arrow/dataset/dataset.h" -#include "arrow/dataset/scanner.h" #include "arrow/record_batch.h" #include "arrow/result.h" #include "arrow/scalar.h" @@ -92,13 +91,15 @@ TEST(SortedMergeNode, Basic) { TEST(SortedMergeNode, TestSortedMergeTwoInputsWithBool) { const int64_t row_count = (16 << 10); // 16k rows per input + // Not a multiple of 8, so that batches start at non-byte-aligned offsets + const int64_t batch_size = 1001; // Create schema with int column A and bool column B auto test_schema = arrow::schema( {arrow::field("col_a", arrow::int32()), arrow::field("col_b", arrow::boolean())}); - // Helper lambda to create table with specific pattern - auto create_test_scanner = [&](int64_t cnt, int offset) -> arrow::Result { + // Helper lambda to create table source with specific pattern + auto make_source = [&](int64_t cnt, int offset) -> arrow::Result { // Create column A (int) - values from offset to offset+cnt-1 arrow::Int32Builder col_a_builder; std::vector col_a_values; @@ -110,10 +111,15 @@ TEST(SortedMergeNode, TestSortedMergeTwoInputsWithBool) { std::shared_ptr col_a_arr; ARROW_RETURN_NOT_OK(col_a_builder.Finish(&col_a_arr)); - // Create column B (bool) - pattern: true if col_a % 5 == 0, false otherwise + // Create column B (bool) - pattern: null if col_a % 7 == 0, otherwise + // true if col_a % 5 == 0, false otherwise arrow::BooleanBuilder col_b_builder; for (int64_t i = 0; i < cnt; ++i) { - int32_t a_value = offset + i; + int32_t a_value = static_cast(offset + i); + if (a_value % 7 == 0) { + ARROW_RETURN_NOT_OK(col_b_builder.AppendNull()); + continue; + } bool b_value = (a_value % 5 == 0); ARROW_RETURN_NOT_OK(col_b_builder.Append(b_value)); } @@ -122,12 +128,12 @@ TEST(SortedMergeNode, TestSortedMergeTwoInputsWithBool) { auto table = arrow::Table::Make(test_schema, {col_a_arr, col_b_arr}); auto table_source = - Declaration("table_source", TableSourceNodeOptions(table, row_count / 16)); + Declaration("table_source", TableSourceNodeOptions(table, batch_size)); return table_source; }; - ASSERT_OK_AND_ASSIGN(auto source1, create_test_scanner(row_count, 0)); - ASSERT_OK_AND_ASSIGN(auto source2, create_test_scanner(row_count, 8192)); + ASSERT_OK_AND_ASSIGN(auto source1, make_source(row_count, 0)); + ASSERT_OK_AND_ASSIGN(auto source2, make_source(row_count, 8192)); // Create sorted merge by column A auto ops = OrderByNodeOptions(compute::Ordering({compute::SortKey("col_a")})); @@ -148,8 +154,6 @@ TEST(SortedMergeNode, TestSortedMergeTwoInputsWithBool) { // Verify sorting and bool values int32_t last_a_value = std::numeric_limits::min(); int64_t total_rows_checked = 0; - int64_t true_count = 0; - int64_t false_count = 0; for (int i = 0; i < col_a->num_chunks(); i++) { auto a_chunk = std::static_pointer_cast(col_a->chunk(i)); @@ -160,10 +164,8 @@ TEST(SortedMergeNode, TestSortedMergeTwoInputsWithBool) { for (int64_t j = 0; j < a_chunk->length(); j++) { ASSERT_FALSE(a_chunk->IsNull(j)) << "Column A should not have null values"; - ASSERT_FALSE(b_chunk->IsNull(j)) << "Column B should not have null values"; int32_t a_value = a_chunk->Value(j); - bool b_value = b_chunk->Value(j); // Verify sorting by column A ASSERT_GE(a_value, last_a_value) @@ -171,18 +173,22 @@ TEST(SortedMergeNode, TestSortedMergeTwoInputsWithBool) { << ": current=" << a_value << ", last=" << last_a_value; last_a_value = a_value; - // Verify bool value correctness: should be true if a_value % 3 == 0 - bool expected_b_value = (a_value % 5 == 0); - ASSERT_EQ(b_value, expected_b_value) - << "Bool value incorrect at chunk " << i << ", row " << j - << ": col_a=" << a_value << ", col_b=" << b_value - << ", expected=" << expected_b_value; - - if (b_value) { - true_count++; - } else { - false_count++; + // Verify bool validity: should be null if a_value % 7 == 0 + bool expected_b_null = (a_value % 7 == 0); + ASSERT_EQ(b_chunk->IsNull(j), expected_b_null) + << "Bool validity incorrect at chunk " << i << ", row " << j + << ": col_a=" << a_value; + + if (!expected_b_null) { + // Verify bool value correctness: should be true if a_value % 5 == 0 + bool b_value = b_chunk->Value(j); + bool expected_b_value = (a_value % 5 == 0); + ASSERT_EQ(b_value, expected_b_value) + << "Bool value incorrect at chunk " << i << ", row " << j + << ": col_a=" << a_value << ", col_b=" << b_value + << ", expected=" << expected_b_value; } + total_rows_checked++; } } @@ -190,7 +196,7 @@ TEST(SortedMergeNode, TestSortedMergeTwoInputsWithBool) { ASSERT_EQ(last_a_value, 24575); ASSERT_EQ(total_rows_checked, row_count * 2) - << "Expected " << row_count << " unique rows after merge"; + << "Expected " << row_count * 2 << " rows after merge"; } } // namespace arrow::acero From f55209b809be8669e75d5d42128c1bfb7f4c7c84 Mon Sep 17 00:00:00 2001 From: Nic Crane Date: Wed, 30 Sep 2026 15:37:33 +0100 Subject: [PATCH 3/4] Update cpp/src/arrow/acero/unmaterialized_table_internal.h Co-authored-by: Rossi Sun --- cpp/src/arrow/acero/unmaterialized_table_internal.h | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/cpp/src/arrow/acero/unmaterialized_table_internal.h b/cpp/src/arrow/acero/unmaterialized_table_internal.h index 2c6b055fc43..17e4da8244f 100644 --- a/cpp/src/arrow/acero/unmaterialized_table_internal.h +++ b/cpp/src/arrow/acero/unmaterialized_table_internal.h @@ -173,11 +173,10 @@ class UnmaterializedCompositeTable { return Status::OK(); } - int64_t array_offset = (source->offset + row) / 8; - int64_t bit_offset = (source->offset + row) % 8; - - builder.UnsafeAppend(arrow::bit_util::GetBit( - source->template GetValues(1, array_offset), bit_offset)); + const int64_t bit_offset = + source->offset + static_cast(row); + builder.UnsafeAppend(bit_util::GetBit( + source->template GetValues(1, 0), bit_offset)); return Status::OK(); } From a559083ca5174488ce8a7f940a7b369cb6837a19 Mon Sep 17 00:00:00 2001 From: Nic Crane Date: Wed, 30 Sep 2026 16:29:10 +0100 Subject: [PATCH 4/4] Fix lint --- cpp/src/arrow/acero/unmaterialized_table_internal.h | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/cpp/src/arrow/acero/unmaterialized_table_internal.h b/cpp/src/arrow/acero/unmaterialized_table_internal.h index 17e4da8244f..335d421500e 100644 --- a/cpp/src/arrow/acero/unmaterialized_table_internal.h +++ b/cpp/src/arrow/acero/unmaterialized_table_internal.h @@ -173,10 +173,9 @@ class UnmaterializedCompositeTable { return Status::OK(); } - const int64_t bit_offset = - source->offset + static_cast(row); - builder.UnsafeAppend(bit_util::GetBit( - source->template GetValues(1, 0), bit_offset)); + const int64_t bit_offset = source->offset + static_cast(row); + builder.UnsafeAppend( + bit_util::GetBit(source->template GetValues(1, 0), bit_offset)); return Status::OK(); }