Skip to content
Merged
4 changes: 4 additions & 0 deletions cpp/src/arrow/compute/exec/options.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -109,6 +109,10 @@ class ARROW_EXPORT ProjectNodeOptions : public ExecNodeOptions {
};

/// \brief Make a node which aggregates input batches, optionally grouped by keys.
///
/// If the keys attribute is a non-empty vector, then each aggregate in `aggregates` is
/// expected to be a HashAggregate function. If the keys attribute is an empty vector,
/// then each aggregate is assumed to be a ScalarAggregate function.
Comment on lines +113 to +115

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

✔️ thank you!

class ARROW_EXPORT AggregateNodeOptions : public ExecNodeOptions {
public:
explicit AggregateNodeOptions(std::vector<Aggregate> aggregates,
Expand Down
30 changes: 30 additions & 0 deletions cpp/src/arrow/compute/exec/plan_test.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -933,6 +933,36 @@ TEST(ExecPlanExecution, SourceGroupedSum) {
}
}

TEST(ExecPlanExecution, SourceMinMaxScalar) {
Comment thread
drin marked this conversation as resolved.
// Regression test for ARROW-16904
for (bool parallel : {false, true}) {
SCOPED_TRACE(parallel ? "parallel/merged" : "serial");

auto input = MakeGroupableBatches(/*multiplicity=*/parallel ? 100 : 1);
auto minmax_opts = std::make_shared<ScalarAggregateOptions>();
auto expected_result = ExecBatch::Make(
{ScalarFromJSON(struct_({field("min", int32()), field("max", int32())}),
R"({"min": -8, "max": 12})")});

ASSERT_OK_AND_ASSIGN(auto plan, ExecPlan::Make());
AsyncGenerator<util::optional<ExecBatch>> sink_gen;

// NOTE: Test `ScalarAggregateNode` by omitting `keys` attribute
ASSERT_OK(Declaration::Sequence(
{{"source",
SourceNodeOptions{input.schema, input.gen(parallel, /*slow=*/false)}},
{"aggregate", AggregateNodeOptions{
/*aggregates=*/{{"min_max", std::move(minmax_opts),
"i32", "min_max"}},
/*keys=*/{}}},
{"sink", SinkNodeOptions{&sink_gen}}})
.AddToPlan(plan.get()));

ASSERT_THAT(StartAndCollect(plan.get(), sink_gen),
Finishes(ResultWith(UnorderedElementsAreArray({*expected_result}))));
}
}

TEST(ExecPlanExecution, NestedSourceFilter) {
for (bool parallel : {false, true}) {
SCOPED_TRACE(parallel ? "parallel/merged" : "serial");
Expand Down
52 changes: 20 additions & 32 deletions cpp/src/arrow/compute/kernels/aggregate_basic_internal.h
Original file line numberDiff line numberDiff line change
Expand Up@@ -440,13 +440,11 @@ struct MinMaxImpl : public ScalarAggregator {
local.has_nulls = !scalar.is_valid;
this->count += scalar.is_valid;

if (local.has_nulls && !options.skip_nulls) {
this->state = local;
return Status::OK();
if (!local.has_nulls || options.skip_nulls) {
local.MergeOne(internal::UnboxScalar<ArrowType>::Unbox(scalar));
}

local.MergeOne(internal::UnboxScalar<ArrowType>::Unbox(scalar));
this->state = local;
this->state += local;
return Status::OK();
}

Expand All@@ -457,19 +455,15 @@ struct MinMaxImpl : public ScalarAggregator {
local.has_nulls = null_count > 0;
this->count += arr.length() - null_count;

if (local.has_nulls && !options.skip_nulls) {
this->state = local;
return Status::OK();
}

if (local.has_nulls) {
local += ConsumeWithNulls(arr);
} else { // All true values
if (!local.has_nulls) {
for (int64_t i = 0; i < arr.length(); i++) {
local.MergeOne(arr.GetView(i));
}
} else if (local.has_nulls && options.skip_nulls) {
local += ConsumeWithNulls(arr);
}
this->state = local;

this->state += local;
return Status::OK();
}

Expand DownExpand Up@@ -585,17 +579,14 @@ struct BooleanMinMaxImpl : public MinMaxImpl<BooleanType, SimdLevel> {

local.has_nulls = null_count > 0;
this->count += valid_count;
if (local.has_nulls && !options.skip_nulls) {
this->state = local;
return Status::OK();
if (!local.has_nulls || options.skip_nulls) {
const auto true_count = arr.true_count();
const auto false_count = valid_count - true_count;
local.max = true_count > 0;
local.min = false_count == 0;
}

const auto true_count = arr.true_count();
const auto false_count = valid_count - true_count;
local.max = true_count > 0;
local.min = false_count == 0;

this->state = local;
this->state += local;
return Status::OK();
}

Expand All@@ -604,17 +595,14 @@ struct BooleanMinMaxImpl : public MinMaxImpl<BooleanType, SimdLevel> {

local.has_nulls = !scalar.is_valid;
this->count += scalar.is_valid;
if (local.has_nulls && !options.skip_nulls) {
this->state = local;
return Status::OK();
if (!local.has_nulls || options.skip_nulls) {
const int true_count = scalar.is_valid && scalar.value;
const int false_count = scalar.is_valid && !scalar.value;
local.max = true_count > 0;
local.min = false_count == 0;
}

const int true_count = scalar.is_valid && scalar.value;
const int false_count = scalar.is_valid && !scalar.value;
local.max = true_count > 0;
local.min = false_count == 0;

this->state = local;
this->state += local;
return Status::OK();
}
};
Expand Down
2 changes: 1 addition & 1 deletion cpp/src/arrow/compute/kernels/aggregate_test.cc
Original file line numberDiff line numberDiff line change
Expand Up@@ -1563,7 +1563,7 @@ TEST_F(TestBooleanMinMaxKernel, Basics) {
TYPED_TEST_SUITE(TestIntegerMinMaxKernel, PhysicalIntegralArrowTypes);
TYPED_TEST(TestIntegerMinMaxKernel, Basics) {
ScalarAggregateOptions options;
std::vector<std::string> chunked_input1 = {"[5, 1, 2, 3, 4]", "[9, 1, null, 3, 4]"};
std::vector<std::string> chunked_input1 = {"[5, 1, 2, 3, 4]", "[9, 8, null, 3, 4]"};
std::vector<std::string> chunked_input2 = {"[5, null, 2, 3, 4]", "[9, 1, 2, 3, 4]"};
std::vector<std::string> chunked_input3 = {"[5, 1, 2, 3, null]", "[9, 1, null, 3, 4]"};
auto item_ty = default_type_instance<TypeParam>();
Expand Down
11 changes: 11 additions & 0 deletions r/tests/testthat/test-dataset.R
Original file line numberDiff line numberDiff line change
Expand Up@@ -618,6 +618,17 @@ test_that("UnionDataset handles InMemoryDatasets", {
expect_equal(actual, expected)
})

test_that("scalar aggregates with many batches (ARROW-16904)", {
tf <- tempfile()
write_parquet(data.frame(x = 1:100), tf, chunk_size = 20)

ds <- open_dataset(tf)
replicate(100, ds %>% summarize(min(x)) %>% pull())

expect_true(all(replicate(100, ds %>% summarize(min(x)) %>% pull()) == 1))
expect_true(all(replicate(100, ds %>% summarize(max(x)) %>% pull()) == 100))
})

test_that("map_batches", {
ds <- open_dataset(dataset_dir, partitioning = "part")

Expand Down