This is an automated email from the ASF dual-hosted git repository.

zeroshade pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-go.git


The following commit(s) were added to refs/heads/main by this push:
     new e9a7361c feat(arrow/compute): scalar aggregate framework with the 
count and sum kernels (#1336)
e9a7361c is described below

commit e9a7361cca27a1049d5046914589424ba9b33b51
Author: singhpratech <[email protected]>
AuthorDate: Fri Sep 25 12:51:08 2026 -0400

    feat(arrow/compute): scalar aggregate framework with the count and sum 
kernels (#1336)
    
    ### Rationale for this change
    
    The compute package has no way to compute a summary value from an array.
    `FuncScalarAgg` exists as
    a function kind but nothing implements it: `exec` has a kernel type for
    scalar and vector kernels
    only, there is no aggregate function type, and `execInternal` returns
    `ErrNotImplemented` for that
    kind. Adding a single aggregate function therefore means adding the
    framework first, which is what
    this pull request does, together with the first two kernels.
    
    The design was discussed on #1296 and the interface follows the changes
    asked for there.
    
    ### What changes are included in this PR?
    
    The framework, in `arrow/compute/exec`:
    
    * `ScalarAggKernel`, following the C++ `ScalarAggregateKernel`: an init
    function that creates a
    state, `Consume` to fold one `ExecSpan` into it, `Merge` to combine two
    states, `Finalize` to
      produce the result, plus an optional `Cleanup` and an `Ordered` flag.
    * `AggregateResult`, a small carrier which holds either an owned
    `scalar.Scalar` or owned
    `arrow.ArrayData`. `exec` cannot import `compute` and so cannot name a
    `Datum`, but an aggregate
    result is not always a scalar, so the carrier keeps the exported
    interface general; the compute
      executor boxes it into a `Datum`.
    * `AggKernel`, the interface the executor consumes, and `MergeAll`, for
    a caller which aggregated
      partitions of the input separately.
    
    The executor, in `arrow/compute`:
    
    * `ScalarAggregateFunction`, `funcImpl[exec.ScalarAggKernel]` with exact
    dispatch, its own
    `SetDefaultOptions`, and `AddKernel`/`AddNewKernel` which reject a
    kernel missing any of the four
      lifecycle functions.
    * `scalarAggExecutor`, reached from `execInternal` for `FuncScalarAgg`.
    It iterates the input in
    spans without promoting scalars to arrays, as the C++
    `ScalarAggExecutor` does, so a kernel sees
    a scalar weighted by the length of the span; it consumes into a single
    state and finalizes once.
    An empty input still reaches finalize, which is what lets `min_count`
    decide the result of an
      aggregation over no values.
    * State cleanup that runs exactly once: when `Execute` returns, whether
    it finalized or failed
    part way, and through the executor's `Clear` on the paths where
    `Execute` never ran, such as a
    failure in `Init`. The result finalize returned owns its own value and
    is unaffected.
    
    The kernels, in `arrow/compute/internal/kernels/aggregate_basic.go` and
    `arrow/compute/scalar_aggregate.go`:
    
    * `count`, with `CountOptions` and the modes `CountOnlyValid`,
    `CountOnlyNull` and `CountAllRows`,
    over any input type. Counting the valid or the null values of a run-end
    encoded, dictionary or
    union array needs a logical null count that an `ArraySpan` does not
    carry, so those return
    `ErrNotImplemented`; `CountAllRows` works for them because it never
    looks at validity.
    * `sum` over int8..int64 (accumulating into int64), uint8..uint64 and
    bool (into uint64),
    float32/float64 (into float64) and the null type (into int64), with the
    C++ semantics for nulls,
    `skip_nulls` and `min_count`, wrap-around on integer overflow, and the
    C++ pairwise summation for
      floating point input.
    * `ScalarAggregateOptions`, `DefaultScalarAggregateOptions`,
    `DefaultCountOptions`, the `Count` and
    `Sum` convenience wrappers, and registration of both option types for
    deserialization.
    
    Aggregates stay out of `exprs`, which rejects any function whose kind is
    not `FuncScalar`
    (`arrow/compute/exprs/exec.go`, in the `*expr.ScalarFunction` case),
    matching C++ where
    aggregations run through Acero rather than through expressions.
    
    Not in this pull request: `mean`, `min_max`, `min`, `max`, `any` and
    `all`, which follow in a
    second one; and `product`, `count_distinct`, `first`/`last`, decimals,
    the statistical aggregates
    and the `hash_*` family, which needs a grouper.
    
    ### Are these changes tested?
    
    Yes.
    
    * Table-driven tests per function over every input type, with and
    without nulls, all null, empty,
    sliced, chunked and scalar input, under the default options,
    `skip_nulls=false`, and `min_count`
    above and below the number of valid values. The expected values were
    taken from the C++
    implementation through pyarrow for the same fixtures and options; all 91
    of them agree exactly.
    * A differential fuzz run of 14,033 cases across the twelve input types,
    plain, sliced, chunked and
    scalar, with every option combination, compared against pyarrow: 0
    mismatches, 0 error
      asymmetries, 0 leaked bytes.
    * Three tests cover the framework rather than the kernels. Every
    aggregate runs at chunk sizes 1 to
    65 and has to agree with the single-span result, which catches a kernel
    that writes a cached null
    count back into the span the executor re-slices. A kernel whose result
    is an array, and one whose
    result is a scalar backed by an allocated buffer, check that the value
    finalize returned outlives
    the cleanup of the state which produced it, and that cleanup runs
    exactly once on success, on
    empty input, on cancellation, and on a consume, finalize or cleanup
    error. And the executor is
    driven directly with a batch of scalars whose logical length is greater
    than one, which the
      `CallFunction` path cannot produce.
    * The pairwise summation is pinned to 100001 doubles from a fixed
    generator; the test also computes
    a naive left-to-right sum of the same fixture and asserts that it
    differs, so the pinned value can
      only be reached with the C++ summation order.
    * `go build ./...`, `go vet ./arrow/compute/...`, `go test
    ./arrow/compute/...` with and without
      `-tags assert`, `go test -race`, and golangci-lint at 0 issues.
    
    ### Are there any user-facing changes?
    
    Yes, all additive. `arrow/compute` gains the `count` and `sum` functions
    in the default registry,
    the `Count` and `Sum` wrappers, `ScalarAggregateOptions`, `CountOptions`
    and their default helpers,
    and `ScalarAggregateFunction`; `arrow/compute/exec` gains
    `ScalarAggKernel`, `AggregateResult`,
    `AggKernel` and `MergeAll`. Nothing existing changes behaviour:
    `execInternal` previously fell through to its default
    branch, `ErrNotImplemented: direct execution of ScalarAggregate`, and no
    function of that kind was
    registered.
---
 arrow/compute/doc.go                              |  17 +-
 arrow/compute/exec.go                             |   8 +
 arrow/compute/exec/aggregate.go                   | 275 +++++++
 arrow/compute/exec/aggregate_test.go              | 182 +++++
 arrow/compute/executor.go                         | 177 ++++
 arrow/compute/expression.go                       |   1 +
 arrow/compute/functions.go                        | 101 ++-
 arrow/compute/internal/kernels/aggregate_basic.go | 622 ++++++++++++++
 arrow/compute/registry.go                         |   1 +
 arrow/compute/scalar_aggregate.go                 | 147 ++++
 arrow/compute/scalar_aggregate_test.go            | 938 ++++++++++++++++++++++
 11 files changed, 2462 insertions(+), 7 deletions(-)

diff --git a/arrow/compute/doc.go b/arrow/compute/doc.go
index 39d9848b..ea0e3e0e 100644
--- a/arrow/compute/doc.go
+++ b/arrow/compute/doc.go
@@ -30,11 +30,18 @@
 // is_in, list_element and the null checks), vector functions (array_filter,
 // array_take, unique, dictionary_encode, cumulative_sum and the run-end
 // encode/decode functions) and the meta functions cast, filter, take, sort
-// and sort_indices that dispatch to them. Scalar aggregate functions (sum,
-// mean, min_max, count, any, all, variance and so on) and hash aggregate
-// functions (the hash_* family used for group-by) are not implemented yet:
-// FuncScalarAgg and FuncHashAgg exist as function kinds, but no function of
-// either kind is registered, and GetFunction returns false for their names.
+// and sort_indices that dispatch to them. It also holds the scalar aggregate
+// functions count and sum, which fold a whole array into a single value
+// through the Init, Consume, Merge and Finalize lifecycle of an
+// exec.ScalarAggKernel. The remaining scalar aggregates (mean, min_max, min,
+// max, any, all, variance and so on) and the hash aggregate functions (the
+// hash_* family used for group-by) are not implemented yet: FuncHashAgg
+// exists as a function kind, but no function of that kind is registered, and
+// GetFunction returns false for their names.
+//
+// Aggregate functions are deliberately not available from the exprs package,
+// which accepts scalar functions only, matching the C++ implementation where
+// aggregations run through Acero rather than through expressions.
 package compute
 
 //go:generate go tool stringer -type=FuncKind -linecomment
diff --git a/arrow/compute/exec.go b/arrow/compute/exec.go
index 710631c8..5a68f11d 100644
--- a/arrow/compute/exec.go
+++ b/arrow/compute/exec.go
@@ -93,6 +93,14 @@ func execInternal(ctx context.Context, fn Function, opts 
FunctionOptions, passed
                        executor.Clear()
                        vectorExecPool.Put(executor.(*vectorExecutor))
                }()
+       case FuncScalarAgg:
+               executor = scalarAggExecPool.Get().(*scalarAggExecutor)
+               // Clear also runs the aggregate state cleanup for the paths 
where
+               // the executor never got to run, and is a no-op once it has run
+               defer func() {
+                       executor.Clear()
+                       scalarAggExecPool.Put(executor.(*scalarAggExecutor))
+               }()
        default:
                return nil, fmt.Errorf("%w: direct execution of %s", 
arrow.ErrNotImplemented, fn.Kind())
        }
diff --git a/arrow/compute/exec/aggregate.go b/arrow/compute/exec/aggregate.go
new file mode 100644
index 00000000..439b7af9
--- /dev/null
+++ b/arrow/compute/exec/aggregate.go
@@ -0,0 +1,275 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+//go:build go1.18
+
+package exec
+
+import (
+       "github.com/apache/arrow-go/v18/arrow"
+       "github.com/apache/arrow-go/v18/arrow/scalar"
+)
+
+// AggregateResult is the single value produced by a scalar aggregate
+// kernel's finalize function.
+//
+// It exists because this package cannot import the compute package and
+// therefore cannot name a Datum, while an aggregate result is not always a
+// scalar: most aggregations produce one (sum, count, min_max as a struct
+// scalar), but some, such as tdigest, produce an array. AggregateResult can
+// hold either one and the executor in the compute package boxes it into the
+// corresponding Datum.
+//
+// # Ownership
+//
+// An AggregateResult owns exactly one reference to the value it holds. A
+// finalize function returns a result that it owns and hands over to its
+// caller, which must eventually either Release it or Take the value out of
+// it and become responsible for that reference. Cleaning up the aggregate
+// state after finalize must not invalidate the returned result: a kernel
+// whose state owns buffers that the result refers to has to retain them for
+// the result (or copy them) before returning it.
+type AggregateResult struct {
+       sc   scalar.Scalar
+       data arrow.ArrayData
+}
+
+// NewScalarResult constructs an aggregate result which owns the reference to
+// the provided scalar.
+func NewScalarResult(sc scalar.Scalar) *AggregateResult {
+       return &AggregateResult{sc: sc}
+}
+
+// NewArrayResult constructs an aggregate result which owns the reference to
+// the provided array data.
+func NewArrayResult(data arrow.ArrayData) *AggregateResult {
+       return &AggregateResult{data: data}
+}
+
+// IsScalar reports whether this result holds a scalar value.
+func (r *AggregateResult) IsScalar() bool { return r.sc != nil }
+
+// IsArray reports whether this result holds array data.
+func (r *AggregateResult) IsArray() bool { return r.data != nil }
+
+// Scalar returns the scalar this result holds, or nil if it holds array
+// data. The reference remains owned by the result.
+func (r *AggregateResult) Scalar() scalar.Scalar { return r.sc }
+
+// ArrayData returns the array data this result holds, or nil if it holds a
+// scalar. The reference remains owned by the result.
+func (r *AggregateResult) ArrayData() arrow.ArrayData { return r.data }
+
+// Type returns the data type of the value this result holds, or nil if it is
+// empty.
+func (r *AggregateResult) Type() arrow.DataType {
+       switch {
+       case r.sc != nil:
+               return r.sc.DataType()
+       case r.data != nil:
+               return r.data.DataType()
+       default:
+               return nil
+       }
+}
+
+// Take hands the owned value over to the caller and empties this result. The
+// caller takes over the reference and is responsible for releasing it. The
+// returned value is a scalar.Scalar, an arrow.ArrayData, or nil if the
+// result was already empty.
+func (r *AggregateResult) Take() any {
+       switch {
+       case r.sc != nil:
+               sc := r.sc
+               r.sc = nil
+               return sc
+       case r.data != nil:
+               data := r.data
+               r.data = nil
+               return data
+       default:
+               return nil
+       }
+}
+
+// Release releases the reference this result owns, if any, and empties it.
+// It is safe to call more than once.
+func (r *AggregateResult) Release() {
+       if rel, ok := r.sc.(scalar.Releasable); ok {
+               rel.Release()
+       }
+       if r.data != nil {
+               r.data.Release()
+       }
+       r.sc, r.data = nil, nil
+}
+
+// ScalarAggConsume folds one span of input into the aggregate state found in
+// ctx.State. Implementations must not modify the span: the executor reuses
+// one ExecSpan for every chunk of the input.
+type ScalarAggConsume = func(ctx *KernelCtx, span *ExecSpan) error
+
+// ScalarAggMerge folds the src state into the dst state.
+//
+// dst is mutated in place and must hold the combined aggregation afterwards.
+// src and dst must not alias, and both must have come from the same kernel's
+// init function. src is left untouched by the merge and remains owned by the
+// caller, which is free to reuse or clean it up; an implementation that wants
+// to keep anything from src past the call has to take its own reference. When
+// the kernel is Ordered, the caller must merge the partitions in the logical
+// order of the input, that is, dst must cover a prefix of the input and src
+// the partition that directly follows it.
+type ScalarAggMerge = func(ctx *KernelCtx, src, dst KernelState) error
+
+// ScalarAggFinalize produces the result of the aggregation from the state in
+// ctx.State. It returns one owned result, see AggregateResult for what the
+// ownership means for the state that produced it.
+type ScalarAggFinalize = func(ctx *KernelCtx) (*AggregateResult, error)
+
+// ScalarAggCleanup releases any resources held by an aggregate state. It is
+// called exactly once for every state that the init function produced,
+// whether the aggregation succeeded, produced no result because the input was
+// empty, was cancelled, or failed in init, consume, merge or finalize. It
+// must not invalidate a result that finalize already returned.
+type ScalarAggCleanup = func(ctx *KernelCtx, state KernelState) error
+
+// AggKernel builds on the base Kernel interface for aggregate execution
+// kernels, which fold input into a state rather than producing an output
+// value per span.
+type AggKernel interface {
+       Kernel
+       Consume(*KernelCtx, *ExecSpan) error
+       Merge(ctx *KernelCtx, src, dst KernelState) error
+       Finalize(*KernelCtx) (*AggregateResult, error)
+       Cleanup(ctx *KernelCtx, state KernelState) error
+       IsOrdered() bool
+}
+
+// A ScalarAggKernel is the kernel implementation for a scalar aggregate
+// function, which computes a single summary value from array input. The four
+// necessary components of an aggregation kernel are the init, consume, merge
+// and finalize functions:
+//
+//   - init (the embedded Kernel's init function): creates a new state.
+//   - consume: folds one ExecSpan into the state in the KernelCtx.
+//   - merge: combines one state with another.
+//   - finalize: produces the result of the aggregation from the state in the
+//     KernelCtx.
+//
+// A fifth, optional, cleanup function releases the resources of a state.
+type ScalarAggKernel struct {
+       kernel
+
+       ConsumeFn  ScalarAggConsume
+       MergeFn    ScalarAggMerge
+       FinalizeFn ScalarAggFinalize
+       CleanupFn  ScalarAggCleanup
+
+       // Ordered indicates that this kernel requires its input in a defined
+       // order, as an aggregation like "first" does. It is the caller of the
+       // kernel which is responsible for passing the data in that order, and
+       // for merging partitions in the logical order of the input. Kernels
+       // whose result does not depend on the order of the input leave this
+       // false.
+       Ordered bool
+}
+
+// NewScalarAggKernel constructs a new kernel for scalar aggregation,
+// constructing a KernelSignature from the provided input and output types.
+//
+// The init function is required: the executor cannot consume without a state.
+// The cleanup function and the Ordered flag are left at their zero values and
+// can be set on the returned kernel.
+func NewScalarAggKernel(in []InputType, out OutputType, init KernelInitFn,
+       consume ScalarAggConsume, merge ScalarAggMerge, finalize 
ScalarAggFinalize) ScalarAggKernel {
+       return NewScalarAggKernelWithSig(&KernelSignature{
+               InputTypes: in,
+               OutType:    out,
+       }, init, consume, merge, finalize)
+}
+
+// NewScalarAggKernelWithSig is a convenience for when you already have a
+// signature to use for constructing a kernel. It is equivalent to passing the
+// components of the signature (input and output types) to NewScalarAggKernel.
+func NewScalarAggKernelWithSig(sig *KernelSignature, init KernelInitFn,
+       consume ScalarAggConsume, merge ScalarAggMerge, finalize 
ScalarAggFinalize) ScalarAggKernel {
+       return ScalarAggKernel{
+               kernel:     kernel{Signature: sig, Init: init},
+               ConsumeFn:  consume,
+               MergeFn:    merge,
+               FinalizeFn: finalize,
+       }
+}
+
+func (s *ScalarAggKernel) Consume(ctx *KernelCtx, span *ExecSpan) error {
+       return s.ConsumeFn(ctx, span)
+}
+
+func (s *ScalarAggKernel) Merge(ctx *KernelCtx, src, dst KernelState) error {
+       return s.MergeFn(ctx, src, dst)
+}
+
+func (s *ScalarAggKernel) Finalize(ctx *KernelCtx) (*AggregateResult, error) {
+       return s.FinalizeFn(ctx)
+}
+
+// Cleanup calls this kernel's cleanup function for the provided state, or
+// does nothing if the kernel does not have one.
+func (s *ScalarAggKernel) Cleanup(ctx *KernelCtx, state KernelState) error {
+       if s.CleanupFn == nil {
+               return nil
+       }
+       return s.CleanupFn(ctx, state)
+}
+
+func (s *ScalarAggKernel) IsOrdered() bool { return s.Ordered }
+
+var _ AggKernel = (*ScalarAggKernel)(nil)
+
+// MergeAll folds every state in states into the first one and returns it,
+// which is the state the caller then finalizes. It is the entry point for a
+// caller which aggregated partitions of the input separately, and mirrors the
+// C++ ScalarAggregateKernel::MergeAll.
+//
+// The states must all have come from this kernel's init function, and, when
+// the kernel is Ordered, must be in the logical order of the input. Every
+// state but the first is cleaned up once it has been merged in, whether or
+// not the merge succeeded; on an error the first state is cleaned up as well
+// and the returned state is nil, so that the caller never has to clean up a
+// state MergeAll was given.
+func MergeAll(kernel AggKernel, ctx *KernelCtx, states []KernelState) 
(KernelState, error) {
+       if len(states) == 0 {
+               return nil, nil
+       }
+
+       dst := states[0]
+       var err error
+       for _, src := range states[1:] {
+               if err == nil {
+                       err = kernel.Merge(ctx, src, dst)
+               }
+               if cleanupErr := kernel.Cleanup(ctx, src); cleanupErr != nil && 
err == nil {
+                       err = cleanupErr
+               }
+       }
+
+       if err != nil {
+               // best effort: the merge error is the one worth reporting
+               _ = kernel.Cleanup(ctx, dst)
+               return nil, err
+       }
+       return dst, nil
+}
diff --git a/arrow/compute/exec/aggregate_test.go 
b/arrow/compute/exec/aggregate_test.go
new file mode 100644
index 00000000..ea208964
--- /dev/null
+++ b/arrow/compute/exec/aggregate_test.go
@@ -0,0 +1,182 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+//go:build go1.18
+
+package exec_test
+
+import (
+       "errors"
+       "testing"
+
+       "github.com/apache/arrow-go/v18/arrow"
+       "github.com/apache/arrow-go/v18/arrow/array"
+       "github.com/apache/arrow-go/v18/arrow/compute/exec"
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/arrow-go/v18/arrow/scalar"
+       "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+)
+
+func TestAggregateResultScalar(t *testing.T) {
+       res := exec.NewScalarResult(scalar.NewInt64Scalar(7))
+       assert.True(t, res.IsScalar())
+       assert.False(t, res.IsArray())
+       assert.Nil(t, res.ArrayData())
+       assert.True(t, arrow.TypeEqual(arrow.PrimitiveTypes.Int64, res.Type()))
+       assert.EqualValues(t, 7, res.Scalar().(*scalar.Int64).Value)
+
+       taken := res.Take()
+       assert.IsType(t, (*scalar.Int64)(nil), taken)
+       assert.False(t, res.IsScalar())
+       assert.Nil(t, res.Take())
+       assert.Nil(t, res.Type())
+
+       // releasing an emptied result is a no-op, and releasing twice is safe
+       res.Release()
+       res.Release()
+}
+
+func TestAggregateResultArrayOwnership(t *testing.T) {
+       mem := memory.NewCheckedAllocator(memory.DefaultAllocator)
+       defer mem.AssertSize(t, 0)
+
+       bldr := array.NewFloat64Builder(mem)
+       bldr.AppendValues([]float64{1, 2, 3}, nil)
+       arr := bldr.NewArray()
+       bldr.Release()
+
+       // the result takes over the single reference held by the array data
+       data := arr.Data()
+       data.Retain()
+       arr.Release()
+
+       res := exec.NewArrayResult(data)
+       assert.True(t, res.IsArray())
+       assert.False(t, res.IsScalar())
+       assert.Nil(t, res.Scalar())
+       assert.True(t, arrow.TypeEqual(arrow.PrimitiveTypes.Float64, 
res.Type()))
+
+       res.Release()
+       assert.False(t, res.IsArray())
+       res.Release()
+}
+
+func TestAggregateResultReleasesBufferBackedScalar(t *testing.T) {
+       mem := memory.NewCheckedAllocator(memory.DefaultAllocator)
+       defer mem.AssertSize(t, 0)
+
+       buf := memory.NewResizableBuffer(mem)
+       buf.Resize(5)
+       copy(buf.Bytes(), "hello")
+
+       res := exec.NewScalarResult(scalar.NewStringScalarFromBuffer(buf))
+       // the scalar took its own reference on the buffer
+       buf.Release()
+
+       assert.Equal(t, "hello", res.Scalar().(*scalar.String).String())
+       res.Release()
+}
+
+// counterKernel is a minimal aggregate kernel used to exercise the framework
+// without going through the compute executor.
+type counterState struct {
+       n        int64
+       mergeErr error
+}
+
+func newCounterKernel(cleanups *int, mergeErr error) exec.ScalarAggKernel {
+       k := exec.NewScalarAggKernel(
+               
[]exec.InputType{exec.NewExactInput(arrow.PrimitiveTypes.Int64)},
+               exec.NewOutputType(arrow.PrimitiveTypes.Int64),
+               func(*exec.KernelCtx, exec.KernelInitArgs) (exec.KernelState, 
error) {
+                       return &counterState{mergeErr: mergeErr}, nil
+               },
+               func(ctx *exec.KernelCtx, span *exec.ExecSpan) error {
+                       ctx.State.(*counterState).n += span.Len
+                       return nil
+               },
+               func(_ *exec.KernelCtx, src, dst exec.KernelState) error {
+                       d := dst.(*counterState)
+                       if d.mergeErr != nil {
+                               return d.mergeErr
+                       }
+                       d.n += src.(*counterState).n
+                       return nil
+               },
+               func(ctx *exec.KernelCtx) (*exec.AggregateResult, error) {
+                       return 
exec.NewScalarResult(scalar.NewInt64Scalar(ctx.State.(*counterState).n)), nil
+               })
+       k.CleanupFn = func(*exec.KernelCtx, exec.KernelState) error {
+               *cleanups++
+               return nil
+       }
+       return k
+}
+
+func TestScalarAggKernelDefaults(t *testing.T) {
+       var cleanups int
+       k := newCounterKernel(&cleanups, nil)
+
+       assert.False(t, k.IsOrdered())
+       assert.NotNil(t, k.GetInitFn())
+       require.NotNil(t, k.GetSig())
+       assert.True(t, 
k.GetSig().MatchesInputs([]arrow.DataType{arrow.PrimitiveTypes.Int64}))
+
+       k.Ordered = true
+       assert.True(t, k.IsOrdered())
+
+       // a kernel without a cleanup function cleans up successfully
+       plain := exec.NewScalarAggKernel(nil, exec.NewOutputType(arrow.Null), 
nil, nil, nil, nil)
+       assert.NoError(t, plain.Cleanup(&exec.KernelCtx{}, nil))
+}
+
+func TestScalarAggMergeAll(t *testing.T) {
+       var cleanups int
+       k := newCounterKernel(&cleanups, nil)
+       ctx := &exec.KernelCtx{Kernel: &k}
+
+       states := []exec.KernelState{
+               &counterState{n: 1}, &counterState{n: 2}, &counterState{n: 3},
+       }
+       merged, err := exec.MergeAll(&k, ctx, states)
+       require.NoError(t, err)
+       assert.EqualValues(t, 6, merged.(*counterState).n)
+       // everything but the state that was merged into has been cleaned up
+       assert.Equal(t, 2, cleanups)
+
+       cleanups = 0
+       merged, err = exec.MergeAll(&k, ctx, nil)
+       assert.NoError(t, err)
+       assert.Nil(t, merged)
+       assert.Zero(t, cleanups)
+}
+
+func TestScalarAggMergeAllCleansUpOnError(t *testing.T) {
+       errMerge := errors.New("merge failed")
+       var cleanups int
+       k := newCounterKernel(&cleanups, errMerge)
+       ctx := &exec.KernelCtx{Kernel: &k}
+
+       states := []exec.KernelState{
+               &counterState{n: 1, mergeErr: errMerge}, &counterState{n: 2}, 
&counterState{n: 3},
+       }
+       merged, err := exec.MergeAll(&k, ctx, states)
+       assert.ErrorIs(t, err, errMerge)
+       assert.Nil(t, merged)
+       // every state given to MergeAll is cleaned up exactly once
+       assert.Equal(t, len(states), cleanups)
+}
diff --git a/arrow/compute/executor.go b/arrow/compute/executor.go
index a804aae9..f666fcae 100644
--- a/arrow/compute/executor.go
+++ b/arrow/compute/executor.go
@@ -918,6 +918,9 @@ var (
        vectorExecPool = sync.Pool{
                New: func() any { return &vectorExecutor{} },
        }
+       scalarAggExecPool = sync.Pool{
+               New: func() any { return &scalarAggExecutor{} },
+       }
 )
 
 func checkCanExecuteChunked(k *exec.VectorKernel) error {
@@ -1220,3 +1223,177 @@ func (v *vectorExecutor) execChunked(batch *ExecBatch, 
out chan<- Datum) error {
        }
        return nil
 }
+
+// NewScalarAggExecutor constructs an executor for scalar aggregate kernels,
+// which consume the whole input into a single kernel state and finalize it
+// once into a single result.
+func NewScalarAggExecutor() KernelExecutor { return &scalarAggExecutor{} }
+
+// scalarAggExecutor drives a scalar aggregate kernel. Unlike the scalar and
+// vector executors it produces exactly one Datum, and it does not preallocate
+// or propagate nulls: the kernel owns the whole result.
+//
+// The state the kernel's init function produced is owned by this executor
+// from Init until cleanup, which happens exactly once, whether the
+// aggregation succeeded, saw no input at all, was cancelled, or failed.
+type scalarAggExecutor struct {
+       ctx     *exec.KernelCtx
+       ectx    ExecCtx
+       kernel  exec.AggKernel
+       outType arrow.DataType
+       cleaned bool
+}
+
+func (a *scalarAggExecutor) Init(ctx *exec.KernelCtx, args 
exec.KernelInitArgs) (err error) {
+       a.ctx, a.cleaned = ctx, false
+       k, ok := args.Kernel.(exec.AggKernel)
+       if !ok {
+               return fmt.Errorf("%w: scalar aggregate execution requires an 
aggregate kernel, got %T",
+                       arrow.ErrInvalid, args.Kernel)
+       }
+       a.kernel = k
+       a.ectx = GetExecCtx(ctx.Ctx)
+       a.outType, err = k.GetSig().OutType.Resolve(ctx, args.Inputs)
+       return
+}
+
+func (a *scalarAggExecutor) Execute(ctx context.Context, batch *ExecBatch, 
data chan<- Datum) (err error) {
+       defer func() {
+               if cleanupErr := a.cleanupState(); cleanupErr != nil && err == 
nil {
+                       err = cleanupErr
+               }
+       }()
+
+       if a.ctx.State == nil {
+               return fmt.Errorf("%w: scalar aggregation requires a non-nil 
kernel state", arrow.ErrInvalid)
+       }
+
+       if err = a.consumeBatch(ctx, batch); err != nil {
+               return
+       }
+
+       // an empty input still reaches finalize, so that the options decide
+       // what an aggregation over no values produces
+       var result *exec.AggregateResult
+       if result, err = a.kernel.Finalize(a.ctx); err != nil {
+               return
+       }
+       if result == nil {
+               return fmt.Errorf("%w: scalar aggregate kernel finalized 
without a result", arrow.ErrInvalid)
+       }
+
+       // the result owns its value; boxing it without owning hands that single
+       // reference over to the Datum, which the caller releases
+       out := NewDatumWithoutOwning(result.Take())
+       select {
+       case <-ctx.Done():
+               out.Release()
+               return context.Cause(ctx)
+       case data <- out:
+               return nil
+       }
+}
+
+// consumeBatch folds every chunk of the batch into the kernel state. Scalars
+// are not promoted to arrays, as they are for scalar kernels: an aggregate
+// kernel handles a scalar input weighted by the length of the span, which is
+// both cheaper and what the C++ implementation does.
+func (a *scalarAggExecutor) consumeBatch(ctx context.Context, batch 
*ExecBatch) error {
+       maxChunkSize := a.ectx.ChunkSize
+       if maxChunkSize <= 0 {
+               maxChunkSize = DefaultMaxChunkSize
+       }
+
+       if checkIfAllScalar(batch) && batch.Len > 1 {
+               // a batch of scalars can carry a logical length greater than 
one,
+               // for instance when only the partition columns of a batch were
+               // projected for an aggregation
+               span := ExecSpanFromBatch(batch)
+               for pos := int64(0); pos < batch.Len; {
+                       span.Len = exec.Min(batch.Len-pos, maxChunkSize)
+                       if err := a.consumeSpan(ctx, span); err != nil {
+                               return err
+                       }
+                       pos += span.Len
+               }
+               return nil
+       }
+
+       _, iter, err := iterateExecSpans(batch, maxChunkSize, false)
+       if err != nil {
+               return err
+       }
+
+       for {
+               span, _, ok := iter()
+               if !ok {
+                       break
+               }
+               if span.Len == 0 {
+                       continue
+               }
+               if err := a.consumeSpan(ctx, &span); err != nil {
+                       return err
+               }
+       }
+       return nil
+}
+
+func (a *scalarAggExecutor) consumeSpan(ctx context.Context, span 
*exec.ExecSpan) error {
+       if err := ctx.Err(); err != nil {
+               return context.Cause(ctx)
+       }
+       return a.kernel.Consume(a.ctx, span)
+}
+
+// WrapResults returns the single result of the aggregation. It reads the
+// channel until it is closed, so that the goroutine which produced the result
+// has finished by the time the state is cleaned up.
+func (a *scalarAggExecutor) WrapResults(ctx context.Context, out <-chan Datum, 
_ bool) Datum {
+       var result Datum
+       for datum := range out {
+               if datum == nil {
+                       continue
+               }
+               if result != nil {
+                       datum.Release()
+                       continue
+               }
+               result = datum
+       }
+
+       if ctx.Err() != nil && result != nil {
+               result.Release()
+               return nil
+       }
+       return result
+}
+
+func (a *scalarAggExecutor) CheckResultType(out Datum) error {
+       typ := out.(ArrayLikeDatum).Type()
+       if typ != nil && !arrow.TypeEqual(a.outType, typ) {
+               return fmt.Errorf("%w: kernel type result mismatch: declared as 
%s, actual is %s",
+                       arrow.ErrType, a.outType, typ)
+       }
+       return nil
+}
+
+func (a *scalarAggExecutor) Clear() {
+       // a no-op if Execute already cleaned up; this covers the paths where it
+       // did not run at all, such as a failure in Init
+       _ = a.cleanupState()
+       a.ctx, a.kernel, a.outType = nil, nil, nil
+       a.cleaned = false
+}
+
+// cleanupState releases the aggregate state exactly once. It does not touch
+// the result which finalize returned: that result owns its own reference.
+func (a *scalarAggExecutor) cleanupState() error {
+       if a.cleaned || a.ctx == nil || a.kernel == nil {
+               return nil
+       }
+       a.cleaned = true
+       state := a.ctx.State
+       a.ctx.State = nil
+       return a.kernel.Cleanup(a.ctx, state)
+}
diff --git a/arrow/compute/expression.go b/arrow/compute/expression.go
index b4c29d0c..7a808280 100644
--- a/arrow/compute/expression.go
+++ b/arrow/compute/expression.go
@@ -580,6 +580,7 @@ var (
                FilterOptions{}, NullOptions{}, StrptimeOptions{}, 
MakeStructOptions{},
                DictionaryEncodeOptions{},
                CumulativeOptions{},
+               ScalarAggregateOptions{}, CountOptions{},
        }
 )
 
diff --git a/arrow/compute/functions.go b/arrow/compute/functions.go
index aad36dac..705deb99 100644
--- a/arrow/compute/functions.go
+++ b/arrow/compute/functions.go
@@ -183,9 +183,10 @@ func (b *baseFunction) checkArity(nargs int) error {
 // generic definitions. It will be extended as other kernel types
 // are defined.
 //
-// Currently only ScalarKernels are allowed to be used.
+// Currently ScalarKernels, VectorKernels and ScalarAggKernels are allowed
+// to be used.
 type kernelType interface {
-       exec.ScalarKernel | exec.VectorKernel
+       exec.ScalarKernel | exec.VectorKernel | exec.ScalarAggKernel
 
        // specifying the Kernel interface here allows us to utilize
        // the methods of the Kernel interface on the generic
@@ -372,6 +373,102 @@ func (f *VectorFunction) Execute(ctx context.Context, 
opts FunctionOptions, args
        return execInternal(ctx, f, opts, -1, args...)
 }
 
+// A ScalarAggregateFunction computes a single summary value from array
+// input, such as the sum or the number of values. Its kernels do not produce
+// one output value per input value: they fold the input into a state which
+// the executor finalizes once, after the whole input has been consumed.
+//
+// Kernels are selected by exact dispatch, as they are for a VectorFunction:
+// there is one kernel per input type and the input is cast to it.
+type ScalarAggregateFunction struct {
+       funcImpl[exec.ScalarAggKernel]
+}
+
+// NewScalarAggregateFunction constructs a new ScalarAggregateFunction object
+// with the passed in name, arity and function doc.
+func NewScalarAggregateFunction(name string, arity Arity, doc FunctionDoc) 
*ScalarAggregateFunction {
+       return &ScalarAggregateFunction{
+               funcImpl: funcImpl[exec.ScalarAggKernel]{
+                       baseFunction: baseFunction{
+                               name:  name,
+                               arity: arity,
+                               doc:   doc,
+                               kind:  FuncScalarAgg,
+                       },
+               },
+       }
+}
+
+func (f *ScalarAggregateFunction) SetDefaultOptions(opts FunctionOptions) {
+       f.defaultOpts = opts
+}
+
+func (f *ScalarAggregateFunction) DispatchExact(vals ...arrow.DataType) 
(exec.Kernel, error) {
+       return f.funcImpl.DispatchExact(vals...)
+}
+
+func (f *ScalarAggregateFunction) DispatchBest(vals ...arrow.DataType) 
(exec.Kernel, error) {
+       return f.DispatchExact(vals...)
+}
+
+// AddNewKernel constructs a new kernel with the provided signature and
+// lifecycle functions and then adds it to the function's list of kernels.
+func (f *ScalarAggregateFunction) AddNewKernel(inTypes []exec.InputType, 
outType exec.OutputType,
+       init exec.KernelInitFn, consume exec.ScalarAggConsume, merge 
exec.ScalarAggMerge,
+       finalize exec.ScalarAggFinalize) error {
+       if err := f.checkArity(len(inTypes)); err != nil {
+               return err
+       }
+
+       if f.arity.IsVarArgs && len(inTypes) != 1 {
+               return fmt.Errorf("%w: varargs signatures must have exactly one 
input type", arrow.ErrInvalid)
+       }
+
+       sig := &exec.KernelSignature{
+               InputTypes: inTypes,
+               OutType:    outType,
+               IsVarArgs:  f.arity.IsVarArgs,
+       }
+
+       return f.AddKernel(exec.NewScalarAggKernelWithSig(sig, init, consume, 
merge, finalize))
+}
+
+// AddKernel adds the provided kernel to the list of kernels this function
+// has. A copy of the kernel is added to the slice of kernels, which means
+// that a given kernel object can be created, added and then reused to add
+// other kernels.
+//
+// Every one of the four lifecycle functions is required: the executor cannot
+// aggregate without a state to fold into, and a kernel without a merge
+// function could not be used by a parallel or partitioned caller.
+func (f *ScalarAggregateFunction) AddKernel(k exec.ScalarAggKernel) error {
+       if err := f.checkArity(len(k.Signature.InputTypes)); err != nil {
+               return err
+       }
+
+       if f.arity.IsVarArgs && !k.Signature.IsVarArgs {
+               return fmt.Errorf("%w: function accepts varargs but kernel 
signature does not", arrow.ErrInvalid)
+       }
+
+       switch {
+       case k.Init == nil:
+               return fmt.Errorf("%w: scalar aggregate kernel requires an init 
function", arrow.ErrInvalid)
+       case k.ConsumeFn == nil:
+               return fmt.Errorf("%w: scalar aggregate kernel requires a 
consume function", arrow.ErrInvalid)
+       case k.MergeFn == nil:
+               return fmt.Errorf("%w: scalar aggregate kernel requires a merge 
function", arrow.ErrInvalid)
+       case k.FinalizeFn == nil:
+               return fmt.Errorf("%w: scalar aggregate kernel requires a 
finalize function", arrow.ErrInvalid)
+       }
+
+       f.kernels = append(f.kernels, k)
+       return nil
+}
+
+func (f *ScalarAggregateFunction) Execute(ctx context.Context, opts 
FunctionOptions, args ...Datum) (Datum, error) {
+       return execInternal(ctx, f, opts, -1, args...)
+}
+
 // MetaFunctionImpl is the signature needed for implementing a MetaFunction
 // which is a function that dispatches to another function instead.
 type MetaFunctionImpl func(context.Context, FunctionOptions, ...Datum) (Datum, 
error)
diff --git a/arrow/compute/internal/kernels/aggregate_basic.go 
b/arrow/compute/internal/kernels/aggregate_basic.go
new file mode 100644
index 00000000..a57a5f12
--- /dev/null
+++ b/arrow/compute/internal/kernels/aggregate_basic.go
@@ -0,0 +1,622 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+//go:build go1.18
+
+package kernels
+
+import (
+       "fmt"
+       "math/bits"
+
+       "github.com/apache/arrow-go/v18/arrow"
+       "github.com/apache/arrow-go/v18/arrow/array"
+       "github.com/apache/arrow-go/v18/arrow/bitutil"
+       "github.com/apache/arrow-go/v18/arrow/compute/exec"
+       "github.com/apache/arrow-go/v18/arrow/scalar"
+       "github.com/apache/arrow-go/v18/internal/bitutils"
+)
+
+// ScalarAggregateOptions controls the general behaviour of the scalar
+// aggregate kernels. The field names and tags are the ones the C++
+// implementation uses so that serialized expressions round-trip.
+//
+// The zero value is not the default: use DefaultScalarAggregateOptions for
+// the behaviour the functions have when they are called without options.
+type ScalarAggregateOptions struct {
+       // SkipNulls, if true, ignores null values. Otherwise, if any value is
+       // null, the aggregation emits null.
+       SkipNulls bool `compute:"skip_nulls"`
+       // MinCount is the number of non-null values which have to be observed
+       // for the aggregation to emit a value rather than null.
+       MinCount uint32 `compute:"min_count"`
+}
+
+func (ScalarAggregateOptions) TypeName() string { return 
"ScalarAggregateOptions" }
+
+// DefaultScalarAggregateOptions returns the options the scalar aggregate
+// functions use when they are called with nil options: null values are
+// skipped and a single non-null value is enough to produce a result.
+func DefaultScalarAggregateOptions() *ScalarAggregateOptions {
+       return &ScalarAggregateOptions{SkipNulls: true, MinCount: 1}
+}
+
+// CountMode is the enum for the values which CountOptions.Mode can take.
+type CountMode int8
+
+const (
+       // CountOnlyValid counts only non-null values.
+       CountOnlyValid CountMode = iota
+       // CountOnlyNull counts only null values.
+       CountOnlyNull
+       // CountAllRows counts both non-null and null values.
+       CountAllRows
+)
+
+// CountOptions controls the behaviour of the count aggregate kernel.
+type CountOptions struct {
+       Mode CountMode `compute:"mode"`
+}
+
+func (CountOptions) TypeName() string { return "CountOptions" }
+
+// DefaultCountOptions returns the options the count function uses when it is
+// called with nil options, which is to count only the non-null values. This
+// is also the zero value of CountOptions.
+func DefaultCountOptions() *CountOptions { return &CountOptions{Mode: 
CountOnlyValid} }
+
+// aggState is implemented by every aggregate state in this package, so that
+// one set of consume/merge/finalize functions can serve all of the kernels,
+// as the C++ ScalarAggregator does.
+type aggState interface {
+       consume(span *exec.ExecSpan) error
+       mergeFrom(src aggState) error
+       finalize() (*exec.AggregateResult, error)
+}
+
+// AggregateConsume is the exec.ScalarAggConsume shared by the aggregate
+// kernels in this package.
+func AggregateConsume(ctx *exec.KernelCtx, span *exec.ExecSpan) error {
+       st, ok := ctx.State.(aggState)
+       if !ok {
+               return fmt.Errorf("%w: aggregate kernel state of unexpected 
type %T", arrow.ErrInvalid, ctx.State)
+       }
+       return st.consume(span)
+}
+
+// AggregateMerge is the exec.ScalarAggMerge shared by the aggregate kernels
+// in this package. It folds src into dst, leaving src untouched.
+func AggregateMerge(_ *exec.KernelCtx, src, dst exec.KernelState) error {
+       dstState, ok := dst.(aggState)
+       if !ok {
+               return fmt.Errorf("%w: aggregate kernel state of unexpected 
type %T", arrow.ErrInvalid, dst)
+       }
+       srcState, ok := src.(aggState)
+       if !ok {
+               return fmt.Errorf("%w: aggregate kernel state of unexpected 
type %T", arrow.ErrInvalid, src)
+       }
+       return dstState.mergeFrom(srcState)
+}
+
+// AggregateFinalize is the exec.ScalarAggFinalize shared by the aggregate
+// kernels in this package.
+func AggregateFinalize(ctx *exec.KernelCtx) (*exec.AggregateResult, error) {
+       st, ok := ctx.State.(aggState)
+       if !ok {
+               return nil, fmt.Errorf("%w: aggregate kernel state of 
unexpected type %T", arrow.ErrInvalid, ctx.State)
+       }
+       return st.finalize()
+}
+
+// spanNullCount returns the number of nulls in the span without writing to
+// it. ArraySpan.UpdateNullCount caches its answer in the span, and the
+// executor re-slices one span for every chunk of the input, so a kernel which
+// updated the count of the span it was handed would leave a stale count
+// behind for the following chunk.
+func spanNullCount(a *exec.ArraySpan) int64 {
+       if a.Type.ID() == arrow.NULL {
+               return a.Len
+       }
+       if nulls := a.Nulls; nulls != array.UnknownNullCount {
+               return nulls
+       }
+       if len(a.Buffers[0].Buf) == 0 {
+               return 0
+       }
+       return a.Len - int64(bitutil.CountSetBits(a.Buffers[0].Buf, 
int(a.Offset), int(a.Len)))
+}
+
+func resolveAggOptions(opts any) (ScalarAggregateOptions, error) {
+       switch o := opts.(type) {
+       case nil:
+               return *DefaultScalarAggregateOptions(), nil
+       case *ScalarAggregateOptions:
+               if o == nil {
+                       return *DefaultScalarAggregateOptions(), nil
+               }
+               return *o, nil
+       case ScalarAggregateOptions:
+               return o, nil
+       }
+       return ScalarAggregateOptions{}, fmt.Errorf("%w: attempted to 
initialize a scalar aggregate kernel with options of type %T",
+               arrow.ErrInvalid, opts)
+}
+
+func resolveCountOptions(opts any) (CountOptions, error) {
+       var out CountOptions
+       switch o := opts.(type) {
+       case nil:
+               out = *DefaultCountOptions()
+       case *CountOptions:
+               if o == nil {
+                       out = *DefaultCountOptions()
+               } else {
+                       out = *o
+               }
+       case CountOptions:
+               out = o
+       default:
+               return out, fmt.Errorf("%w: attempted to initialize the count 
kernel with options of type %T",
+                       arrow.ErrInvalid, opts)
+       }
+
+       switch out.Mode {
+       case CountOnlyValid, CountOnlyNull, CountAllRows:
+       default:
+               return out, fmt.Errorf("%w: invalid count mode %d", 
arrow.ErrInvalid, out.Mode)
+       }
+       return out, nil
+}
+
+// ----------------------------------------------------------------------
+// count
+
+type countImpl struct {
+       opts     CountOptions
+       nulls    int64
+       nonNulls int64
+}
+
+func (c *countImpl) consume(span *exec.ExecSpan) error {
+       v := &span.Values[0]
+       switch {
+       case c.opts.Mode == CountAllRows:
+               // ALL never has to look at the validity bitmap.
+               c.nonNulls += span.Len
+       case v.IsArray():
+               nulls := spanNullCount(&v.Array)
+               c.nulls += nulls
+               c.nonNulls += v.Array.Len - nulls
+       case v.Scalar.IsValid():
+               c.nonNulls += span.Len
+       default:
+               c.nulls += span.Len
+       }
+       return nil
+}
+
+func (c *countImpl) mergeFrom(src aggState) error {
+       other, ok := src.(*countImpl)
+       if !ok {
+               return fmt.Errorf("%w: cannot merge %T into a count state", 
arrow.ErrInvalid, src)
+       }
+       c.nulls += other.nulls
+       c.nonNulls += other.nonNulls
+       return nil
+}
+
+func (c *countImpl) finalize() (*exec.AggregateResult, error) {
+       // ALL is equivalent to ONLY_VALID here because consume counted every
+       // row as non-null rather than computing a null count it would not use.
+       if c.opts.Mode == CountOnlyNull {
+               return exec.NewScalarResult(scalar.NewInt64Scalar(c.nulls)), nil
+       }
+       return exec.NewScalarResult(scalar.NewInt64Scalar(c.nonNulls)), nil
+}
+
+// CountInit is the exec.KernelInitFn for the count kernel.
+func CountInit(_ *exec.KernelCtx, args exec.KernelInitArgs) (exec.KernelState, 
error) {
+       opts, err := resolveCountOptions(args.Options)
+       if err != nil {
+               return nil, err
+       }
+
+       // The logical null count of these types is not the number of unset bits
+       // in the validity bitmap of the top-level array, so counting them needs
+       // machinery this kernel does not have yet.
+       switch args.Inputs[0].ID() {
+       case arrow.RUN_END_ENCODED, arrow.DICTIONARY, arrow.SPARSE_UNION, 
arrow.DENSE_UNION:
+               if opts.Mode != CountAllRows {
+                       return nil, fmt.Errorf("%w: count of %s requires a 
logical null count",
+                               arrow.ErrNotImplemented, args.Inputs[0])
+               }
+       }
+
+       return &countImpl{opts: opts}, nil
+}
+
+// ----------------------------------------------------------------------
+// sum
+
+// sumBase holds the parts of a sum state which do not depend on the type
+// being summed: how many values were seen, whether any null was seen, and the
+// options which turn the two into a null or a value at finalize time.
+type sumBase struct {
+       opts          ScalarAggregateOptions
+       count         uint64
+       nullsObserved bool
+}
+
+// observe folds the value and null counts of one span into the state. It
+// returns the number of non-null values in the span and whether their values
+// still have to be accumulated: once a null has been seen with SkipNulls
+// false, the result is null whatever the values are.
+func (s *sumBase) observe(span *exec.ExecSpan) (valid int64, accumulate bool) {
+       v := &span.Values[0]
+       if v.IsArray() {
+               nulls := spanNullCount(&v.Array)
+               valid = v.Array.Len - nulls
+               s.nullsObserved = s.nullsObserved || nulls > 0
+       } else if v.Scalar.IsValid() {
+               valid = span.Len
+       } else {
+               s.nullsObserved = true
+       }
+       s.count += uint64(valid)
+       return valid, valid > 0 && (s.opts.SkipNulls || !s.nullsObserved)
+}
+
+func (s *sumBase) mergeBase(other *sumBase) {
+       s.count += other.count
+       s.nullsObserved = s.nullsObserved || other.nullsObserved
+}
+
+// emitNull reports whether finalize has to produce a null rather than the
+// accumulated value.
+func (s *sumBase) emitNull() bool {
+       return (!s.opts.SkipNulls && s.nullsObserved) || s.count < 
uint64(s.opts.MinCount)
+}
+
+// intSumImpl sums the signed integer types into an int64. Go defines the
+// wrap-around, so an overflowing sum wraps rather than being undefined as it
+// is in C++, but the value produced is the same one C++ produces in practice.
+type intSumImpl[T int8 | int16 | int32 | int64] struct {
+       sumBase
+       sum int64
+}
+
+func (s *intSumImpl[T]) consume(span *exec.ExecSpan) error {
+       if _, accumulate := s.observe(span); !accumulate {
+               return nil
+       }
+       v := &span.Values[0]
+       if v.IsArray() {
+               visitValues(&v.Array, func(vals []T, pos, length int64) {
+                       for i := int64(0); i < length; i++ {
+                               s.sum += int64(vals[pos+i])
+                       }
+               })
+       } else {
+               s.sum += 
int64(UnboxScalar[T](v.Scalar.(scalar.PrimitiveScalar))) * span.Len
+       }
+       return nil
+}
+
+func (s *intSumImpl[T]) mergeFrom(src aggState) error {
+       other, ok := src.(*intSumImpl[T])
+       if !ok {
+               return fmt.Errorf("%w: cannot merge %T into a sum state", 
arrow.ErrInvalid, src)
+       }
+       s.mergeBase(&other.sumBase)
+       s.sum += other.sum
+       return nil
+}
+
+func (s *intSumImpl[T]) finalize() (*exec.AggregateResult, error) {
+       if s.emitNull() {
+               return 
exec.NewScalarResult(scalar.MakeNullScalar(arrow.PrimitiveTypes.Int64)), nil
+       }
+       return exec.NewScalarResult(scalar.NewInt64Scalar(s.sum)), nil
+}
+
+// uintSumImpl sums the unsigned integer types into a uint64.
+type uintSumImpl[T uint8 | uint16 | uint32 | uint64] struct {
+       sumBase
+       sum uint64
+}
+
+func (s *uintSumImpl[T]) consume(span *exec.ExecSpan) error {
+       if _, accumulate := s.observe(span); !accumulate {
+               return nil
+       }
+       v := &span.Values[0]
+       if v.IsArray() {
+               visitValues(&v.Array, func(vals []T, pos, length int64) {
+                       for i := int64(0); i < length; i++ {
+                               s.sum += uint64(vals[pos+i])
+                       }
+               })
+       } else {
+               s.sum += 
uint64(UnboxScalar[T](v.Scalar.(scalar.PrimitiveScalar))) * uint64(span.Len)
+       }
+       return nil
+}
+
+func (s *uintSumImpl[T]) mergeFrom(src aggState) error {
+       other, ok := src.(*uintSumImpl[T])
+       if !ok {
+               return fmt.Errorf("%w: cannot merge %T into a sum state", 
arrow.ErrInvalid, src)
+       }
+       s.mergeBase(&other.sumBase)
+       s.sum += other.sum
+       return nil
+}
+
+func (s *uintSumImpl[T]) finalize() (*exec.AggregateResult, error) {
+       if s.emitNull() {
+               return 
exec.NewScalarResult(scalar.MakeNullScalar(arrow.PrimitiveTypes.Uint64)), nil
+       }
+       return exec.NewScalarResult(scalar.NewUint64Scalar(s.sum)), nil
+}
+
+// floatSumImpl sums the floating point types into a float64 using the same
+// pairwise summation the C++ implementation uses, so that the result is the
+// same one down to the last bit.
+type floatSumImpl[T float32 | float64] struct {
+       sumBase
+       sum float64
+}
+
+func (s *floatSumImpl[T]) consume(span *exec.ExecSpan) error {
+       valid, accumulate := s.observe(span)
+       if !accumulate {
+               return nil
+       }
+       v := &span.Values[0]
+       if v.IsArray() {
+               s.sum += pairwiseSum[T](&v.Array, valid)
+       } else {
+               // C++ multiplies the unboxed value by the span length in the 
value's
+               // own type before widening it to the accumulator.
+               val := UnboxScalar[T](v.Scalar.(scalar.PrimitiveScalar))
+               s.sum += float64(val * T(span.Len))
+       }
+       return nil
+}
+
+func (s *floatSumImpl[T]) mergeFrom(src aggState) error {
+       other, ok := src.(*floatSumImpl[T])
+       if !ok {
+               return fmt.Errorf("%w: cannot merge %T into a sum state", 
arrow.ErrInvalid, src)
+       }
+       s.mergeBase(&other.sumBase)
+       s.sum += other.sum
+       return nil
+}
+
+func (s *floatSumImpl[T]) finalize() (*exec.AggregateResult, error) {
+       if s.emitNull() {
+               return 
exec.NewScalarResult(scalar.MakeNullScalar(arrow.PrimitiveTypes.Float64)), nil
+       }
+       return exec.NewScalarResult(scalar.NewFloat64Scalar(s.sum)), nil
+}
+
+// boolSumImpl sums boolean input as the number of true values, into a uint64.
+type boolSumImpl struct {
+       sumBase
+       sum uint64
+}
+
+func (s *boolSumImpl) consume(span *exec.ExecSpan) error {
+       if _, accumulate := s.observe(span); !accumulate {
+               return nil
+       }
+       v := &span.Values[0]
+       if v.IsArray() {
+               s.sum += uint64(trueCount(&v.Array))
+       } else if v.Scalar.(*scalar.Boolean).Value {
+               s.sum += uint64(span.Len)
+       }
+       return nil
+}
+
+func (s *boolSumImpl) mergeFrom(src aggState) error {
+       other, ok := src.(*boolSumImpl)
+       if !ok {
+               return fmt.Errorf("%w: cannot merge %T into a sum state", 
arrow.ErrInvalid, src)
+       }
+       s.mergeBase(&other.sumBase)
+       s.sum += other.sum
+       return nil
+}
+
+func (s *boolSumImpl) finalize() (*exec.AggregateResult, error) {
+       if s.emitNull() {
+               return 
exec.NewScalarResult(scalar.MakeNullScalar(arrow.PrimitiveTypes.Uint64)), nil
+       }
+       return exec.NewScalarResult(scalar.NewUint64Scalar(s.sum)), nil
+}
+
+// nullSumImpl sums null-typed input, whose every value is null, into an
+// int64. It reproduces the C++ NullSumImpl: the only thing it has to track is
+// whether anything at all was consumed.
+type nullSumImpl struct {
+       opts    ScalarAggregateOptions
+       isEmpty bool
+}
+
+func (s *nullSumImpl) consume(span *exec.ExecSpan) error {
+       v := &span.Values[0]
+       if v.IsScalar() || spanNullCount(&v.Array) > 0 {
+               s.isEmpty = false
+       }
+       return nil
+}
+
+func (s *nullSumImpl) mergeFrom(src aggState) error {
+       other, ok := src.(*nullSumImpl)
+       if !ok {
+               return fmt.Errorf("%w: cannot merge %T into a sum state", 
arrow.ErrInvalid, src)
+       }
+       s.isEmpty = s.isEmpty && other.isEmpty
+       return nil
+}
+
+func (s *nullSumImpl) finalize() (*exec.AggregateResult, error) {
+       if (s.opts.SkipNulls || s.isEmpty) && s.opts.MinCount == 0 {
+               return exec.NewScalarResult(scalar.NewInt64Scalar(0)), nil
+       }
+       return 
exec.NewScalarResult(scalar.MakeNullScalar(arrow.PrimitiveTypes.Int64)), nil
+}
+
+// SumInit is the exec.KernelInitFn for every kernel of the sum function; it
+// picks the accumulator from the input type the way the C++ SumLikeInit
+// visitor does.
+func SumInit(_ *exec.KernelCtx, args exec.KernelInitArgs) (exec.KernelState, 
error) {
+       opts, err := resolveAggOptions(args.Options)
+       if err != nil {
+               return nil, err
+       }
+
+       base := sumBase{opts: opts}
+       switch args.Inputs[0].ID() {
+       case arrow.INT8:
+               return &intSumImpl[int8]{sumBase: base}, nil
+       case arrow.INT16:
+               return &intSumImpl[int16]{sumBase: base}, nil
+       case arrow.INT32:
+               return &intSumImpl[int32]{sumBase: base}, nil
+       case arrow.INT64:
+               return &intSumImpl[int64]{sumBase: base}, nil
+       case arrow.UINT8:
+               return &uintSumImpl[uint8]{sumBase: base}, nil
+       case arrow.UINT16:
+               return &uintSumImpl[uint16]{sumBase: base}, nil
+       case arrow.UINT32:
+               return &uintSumImpl[uint32]{sumBase: base}, nil
+       case arrow.UINT64:
+               return &uintSumImpl[uint64]{sumBase: base}, nil
+       case arrow.FLOAT32:
+               return &floatSumImpl[float32]{sumBase: base}, nil
+       case arrow.FLOAT64:
+               return &floatSumImpl[float64]{sumBase: base}, nil
+       case arrow.BOOL:
+               return &boolSumImpl{sumBase: base}, nil
+       case arrow.NULL:
+               return &nullSumImpl{opts: opts, isEmpty: true}, nil
+       }
+       return nil, fmt.Errorf("%w: no sum implemented for %s", 
arrow.ErrNotImplemented, args.Inputs[0])
+}
+
+// ----------------------------------------------------------------------
+// value iteration helpers
+
+// visitValues calls fn once for every run of valid values in the span,
+// passing the typed values of the span along with the position and length of
+// the run relative to the start of the span.
+func visitValues[T arrow.FixedWidthType](a *exec.ArraySpan, fn func(vals []T, 
pos, length int64)) {
+       vals := exec.GetSpanValues[T](a, 1)
+       bitutils.VisitSetBitRunsNoErr(a.Buffers[0].Buf, a.Offset, a.Len, 
func(pos, length int64) {
+               fn(vals, pos, length)
+       })
+}
+
+// trueCount returns the number of values in the boolean span which are both
+// valid and true.
+func trueCount(a *exec.ArraySpan) int64 {
+       values := a.Buffers[1].Buf
+       if len(values) == 0 {
+               return 0
+       }
+       if len(a.Buffers[0].Buf) == 0 {
+               return int64(bitutil.CountSetBits(values, int(a.Offset), 
int(a.Len)))
+       }
+
+       var count int64
+       bitutils.VisitSetBitRunsNoErr(a.Buffers[0].Buf, a.Offset, a.Len, 
func(pos, length int64) {
+               count += int64(bitutil.CountSetBits(values, int(a.Offset+pos), 
int(length)))
+       })
+       return count
+}
+
+// pairwiseSum is a port of the non-recursive pairwise summation the C++
+// implementation uses for floating point input, so that the two produce
+// bit-for-bit identical results.
+//
+// https://en.wikipedia.org/wiki/Pairwise_summation
+func pairwiseSum[T float32 | float64](a *exec.ArraySpan, dataSize int64) 
float64 {
+       if dataSize == 0 {
+               return 0
+       }
+
+       // number of inputs to accumulate before merging with another block
+       const blockSize = 16 // same as numpy
+       // levels (tree depth) = ceil(log2(len)) + 1, a bit larger than 
necessary
+       levels := bits.Len64(uint64(dataSize-1)) + 1
+       // temporary summation per level
+       sum := make([]float64, levels+1)
+       // whether two summations are ready and should be reduced to the level
+       // above; one bit per level, bit 0 is level 0, and so on
+       var mask uint64
+       // level of the root node holding the final summation
+       var rootLevel int
+
+       // reduce the summation of one block (which may be smaller than 
blockSize)
+       // from a leaf node, continuing to the level above whenever two 
summations
+       // are ready for a non-leaf node
+       reduce := func(blockSum float64) {
+               curLevel, curLevelMask := 0, uint64(1)
+               sum[curLevel] += blockSum
+               mask ^= curLevelMask
+               for mask&curLevelMask == 0 {
+                       blockSum = sum[curLevel]
+                       sum[curLevel] = 0
+                       curLevel++
+                       curLevelMask <<= 1
+                       sum[curLevel] += blockSum
+                       mask ^= curLevelMask
+               }
+               if curLevel > rootLevel {
+                       rootLevel = curLevel
+               }
+       }
+
+       visitValues[T](a, func(vals []T, pos, length int64) {
+               v := vals[pos:]
+               blocks, remains := length/blockSize, length%blockSize
+               for i := int64(0); i < blocks; i++ {
+                       var blockSum float64
+                       for j := 0; j < blockSize; j++ {
+                               blockSum += float64(v[j])
+                       }
+                       reduce(blockSum)
+                       v = v[blockSize:]
+               }
+               if remains > 0 {
+                       var blockSum float64
+                       for i := int64(0); i < remains; i++ {
+                               blockSum += float64(v[i])
+                       }
+                       reduce(blockSum)
+               }
+       })
+
+       // reduce the intermediate summations from all of the non-leaf nodes
+       for i := 1; i <= rootLevel; i++ {
+               sum[i] += sum[i-1]
+       }
+       return sum[rootLevel]
+}
diff --git a/arrow/compute/registry.go b/arrow/compute/registry.go
index 0cc11025..03d4df81 100644
--- a/arrow/compute/registry.go
+++ b/arrow/compute/registry.go
@@ -58,6 +58,7 @@ func GetFunctionRegistry() FunctionRegistry {
                RegisterVectorCumulative(registry)
                RegisterVectorRunEndFuncs(registry)
                RegisterScalarSetLookup(registry)
+               RegisterScalarAggregates(registry)
        })
        return registry
 }
diff --git a/arrow/compute/scalar_aggregate.go 
b/arrow/compute/scalar_aggregate.go
new file mode 100644
index 00000000..73913526
--- /dev/null
+++ b/arrow/compute/scalar_aggregate.go
@@ -0,0 +1,147 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+//go:build go1.18
+
+package compute
+
+import (
+       "context"
+
+       "github.com/apache/arrow-go/v18/arrow"
+       "github.com/apache/arrow-go/v18/arrow/compute/exec"
+       "github.com/apache/arrow-go/v18/arrow/compute/internal/kernels"
+)
+
+type (
+       // ScalarAggregateOptions controls the general behaviour of the scalar
+       // aggregate functions: whether nulls are skipped and how many values
+       // have to be seen for a result to be produced rather than a null.
+       //
+       // The zero value means SkipNulls false and MinCount zero, which is not
+       // the default behaviour of the functions; see
+       // DefaultScalarAggregateOptions.
+       ScalarAggregateOptions = kernels.ScalarAggregateOptions
+       // CountOptions controls which values the count function counts.
+       CountOptions = kernels.CountOptions
+       // CountMode is the enum of the values CountOptions.Mode can take.
+       CountMode = kernels.CountMode
+)
+
+const (
+       // CountOnlyValid counts only non-null values.
+       CountOnlyValid = kernels.CountOnlyValid
+       // CountOnlyNull counts only null values.
+       CountOnlyNull = kernels.CountOnlyNull
+       // CountAllRows counts both non-null and null values.
+       CountAllRows = kernels.CountAllRows
+)
+
+// DefaultScalarAggregateOptions returns the options which the scalar
+// aggregate functions use when they are called without options: nulls are
+// skipped and one non-null value is enough to produce a result.
+func DefaultScalarAggregateOptions() *ScalarAggregateOptions {
+       return kernels.DefaultScalarAggregateOptions()
+}
+
+// DefaultCountOptions returns the options which the count function uses when
+// it is called without options, which is to count only non-null values.
+func DefaultCountOptions() *CountOptions { return 
kernels.DefaultCountOptions() }
+
+var (
+       countDoc = FunctionDoc{
+               Summary: "Count the number of values",
+               Description: "By default, only non-null values are counted.\n" +
+                       "This can be changed through the mode option.",
+               ArgNames:    []string{"array"},
+               OptionsType: "CountOptions",
+       }
+
+       sumDoc = FunctionDoc{
+               Summary: "Compute the sum of a numeric array",
+               Description: "Null values are ignored by default. If the 
skip_nulls\n" +
+                       "option is set to false, then a null is emitted as soon 
as one of\n" +
+                       "the values is null. A null is also emitted when fewer 
than\n" +
+                       "min_count non-null values are seen, which by default 
means that\n" +
+                       "an empty or an all-null input produces a null.",
+               ArgNames:    []string{"array"},
+               OptionsType: "ScalarAggregateOptions",
+       }
+)
+
+// sumKernelTypes lists the input types the sum function accepts along with
+// the type it accumulates them into, following the C++ implementation: the
+// signed integers widen to int64, the unsigned ones and booleans to uint64,
+// and the floats to float64. Null input sums as int64 so that a sum over a
+// column of unknown type still has a numeric type.
+var sumKernelTypes = []struct {
+       in  arrow.DataType
+       out arrow.DataType
+}{
+       {arrow.PrimitiveTypes.Int8, arrow.PrimitiveTypes.Int64},
+       {arrow.PrimitiveTypes.Int16, arrow.PrimitiveTypes.Int64},
+       {arrow.PrimitiveTypes.Int32, arrow.PrimitiveTypes.Int64},
+       {arrow.PrimitiveTypes.Int64, arrow.PrimitiveTypes.Int64},
+       {arrow.PrimitiveTypes.Uint8, arrow.PrimitiveTypes.Uint64},
+       {arrow.PrimitiveTypes.Uint16, arrow.PrimitiveTypes.Uint64},
+       {arrow.PrimitiveTypes.Uint32, arrow.PrimitiveTypes.Uint64},
+       {arrow.PrimitiveTypes.Uint64, arrow.PrimitiveTypes.Uint64},
+       {arrow.PrimitiveTypes.Float32, arrow.PrimitiveTypes.Float64},
+       {arrow.PrimitiveTypes.Float64, arrow.PrimitiveTypes.Float64},
+       {arrow.FixedWidthTypes.Boolean, arrow.PrimitiveTypes.Uint64},
+       {arrow.Null, arrow.PrimitiveTypes.Int64},
+}
+
+// RegisterScalarAggregates registers the scalar aggregate functions, which
+// compute a single summary value from array input.
+func RegisterScalarAggregates(reg FunctionRegistry) {
+       countFn := NewScalarAggregateFunction("count", Unary(), countDoc)
+       countFn.SetDefaultOptions(DefaultCountOptions())
+       // count is registered for any input type: it only ever looks at the
+       // validity of the values, never at the values themselves. CountInit
+       // rejects the modes that need a logical null count for the types whose
+       // validity bitmap does not carry it.
+       if err := countFn.AddNewKernel([]exec.InputType{{}}, 
exec.NewOutputType(arrow.PrimitiveTypes.Int64),
+               kernels.CountInit, kernels.AggregateConsume, 
kernels.AggregateMerge, kernels.AggregateFinalize); err != nil {
+               panic(err)
+       }
+       reg.AddFunction(countFn, false)
+
+       sumFn := NewScalarAggregateFunction("sum", Unary(), sumDoc)
+       sumFn.SetDefaultOptions(DefaultScalarAggregateOptions())
+       for _, kt := range sumKernelTypes {
+               if err := 
sumFn.AddNewKernel([]exec.InputType{exec.NewExactInput(kt.in)}, 
exec.NewOutputType(kt.out),
+                       kernels.SumInit, kernels.AggregateConsume, 
kernels.AggregateMerge, kernels.AggregateFinalize); err != nil {
+                       panic(err)
+               }
+       }
+       reg.AddFunction(sumFn, false)
+}
+
+// Count returns the number of values in the input, counting either the
+// non-null values, the null values or every row depending on the mode in the
+// provided options.
+func Count(ctx context.Context, opts CountOptions, value Datum) (Datum, error) 
{
+       return CallFunction(ctx, "count", &opts, value)
+}
+
+// Sum returns the sum of the values in the input, as an int64 for signed
+// integer and null input, a uint64 for unsigned integer and boolean input and
+// a float64 for floating point input. The sum of an integer input which does
+// not fit its output type wraps around.
+func Sum(ctx context.Context, opts ScalarAggregateOptions, value Datum) 
(Datum, error) {
+       return CallFunction(ctx, "sum", &opts, value)
+}
diff --git a/arrow/compute/scalar_aggregate_test.go 
b/arrow/compute/scalar_aggregate_test.go
new file mode 100644
index 00000000..b22719a5
--- /dev/null
+++ b/arrow/compute/scalar_aggregate_test.go
@@ -0,0 +1,938 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+//go:build go1.18
+
+package compute_test
+
+import (
+       "context"
+       "errors"
+       "fmt"
+       "math"
+       "strings"
+       "testing"
+
+       "github.com/apache/arrow-go/v18/arrow"
+       "github.com/apache/arrow-go/v18/arrow/array"
+       "github.com/apache/arrow-go/v18/arrow/compute"
+       "github.com/apache/arrow-go/v18/arrow/compute/exec"
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/arrow-go/v18/arrow/scalar"
+       "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+)
+
+// The expected values in this file were produced by the C++ implementation
+// through pyarrow for the same fixtures and options.
+
+func aggArray(t *testing.T, mem memory.Allocator, dt arrow.DataType, jsonVals 
string) arrow.Array {
+       t.Helper()
+       arr, _, err := array.FromJSON(mem, dt, strings.NewReader(jsonVals))
+       require.NoError(t, err)
+       return arr
+}
+
+// callAgg calls the named aggregate function and returns the scalar it
+// produced. Aggregate results are scalars for every kernel in this file.
+func callAgg(t *testing.T, ctx context.Context, fname string, opts 
compute.FunctionOptions, d compute.Datum) scalar.Scalar {
+       t.Helper()
+       res, err := compute.CallFunction(ctx, fname, opts, d)
+       require.NoError(t, err)
+       defer res.Release()
+
+       sd, ok := res.(*compute.ScalarDatum)
+       require.Truef(t, ok, "expected a scalar datum, got %T", res)
+       return sd.Value
+}
+
+func aggOf(t *testing.T, ctx context.Context, fname string, opts 
compute.FunctionOptions, mem memory.Allocator, dt arrow.DataType, jsonVals 
string) scalar.Scalar {
+       t.Helper()
+       arr := aggArray(t, mem, dt, jsonVals)
+       defer arr.Release()
+       return callAgg(t, ctx, fname, opts, compute.NewDatumWithoutOwning(arr))
+}
+
+func assertNullResult(t *testing.T, sc scalar.Scalar, dt arrow.DataType) {
+       t.Helper()
+       assert.Falsef(t, sc.IsValid(), "expected a null result, got %s", sc)
+       assert.Truef(t, arrow.TypeEqual(dt, sc.DataType()), "expected type %s, 
got %s", dt, sc.DataType())
+}
+
+func assertInt64Result(t *testing.T, sc scalar.Scalar, want int64) {
+       t.Helper()
+       v, ok := sc.(*scalar.Int64)
+       require.Truef(t, ok, "expected an int64 scalar, got %T", sc)
+       require.True(t, v.IsValid(), "expected a valid result")
+       assert.Equal(t, want, v.Value)
+}
+
+func assertUint64Result(t *testing.T, sc scalar.Scalar, want uint64) {
+       t.Helper()
+       v, ok := sc.(*scalar.Uint64)
+       require.Truef(t, ok, "expected a uint64 scalar, got %T", sc)
+       require.True(t, v.IsValid(), "expected a valid result")
+       assert.Equal(t, want, v.Value)
+}
+
+func assertFloat64Result(t *testing.T, sc scalar.Scalar, want float64) {
+       t.Helper()
+       v, ok := sc.(*scalar.Float64)
+       require.Truef(t, ok, "expected a float64 scalar, got %T", sc)
+       require.True(t, v.IsValid(), "expected a valid result")
+       assert.Equal(t, want, v.Value)
+}
+
+func aggTestContext(t *testing.T) (context.Context, *memory.CheckedAllocator) {
+       t.Helper()
+       mem := memory.NewCheckedAllocator(memory.DefaultAllocator)
+       t.Cleanup(func() { mem.AssertSize(t, 0) })
+       return compute.WithAllocator(context.Background(), mem), mem
+}
+
+func skipNulls(skip bool, minCount uint32) *compute.ScalarAggregateOptions {
+       return &compute.ScalarAggregateOptions{SkipNulls: skip, MinCount: 
minCount}
+}
+
+func TestScalarAggregateRegistered(t *testing.T) {
+       reg := compute.GetFunctionRegistry()
+       for _, name := range []string{"count", "sum"} {
+               fn, ok := reg.GetFunction(name)
+               require.Truef(t, ok, "function %s is not registered", name)
+               assert.Equal(t, compute.FuncScalarAgg, fn.Kind())
+               assert.Equal(t, compute.Unary(), fn.Arity())
+               assert.NoError(t, fn.Validate())
+               assert.NotNil(t, fn.DefaultOptions())
+               assert.Greater(t, fn.NumKernels(), 0)
+       }
+
+       // exact dispatch: sum has no kernel for a type it does not accept
+       sumFn, _ := reg.GetFunction("sum")
+       _, err := sumFn.DispatchBest(arrow.BinaryTypes.String)
+       assert.ErrorIs(t, err, arrow.ErrNotImplemented)
+       _, err = sumFn.DispatchBest(arrow.FixedWidthTypes.Float16)
+       assert.ErrorIs(t, err, arrow.ErrNotImplemented)
+       k, err := sumFn.DispatchBest(arrow.PrimitiveTypes.Int32)
+       require.NoError(t, err)
+       assert.IsType(t, (*exec.ScalarAggKernel)(nil), k)
+
+       // count takes any input type
+       countFn, _ := reg.GetFunction("count")
+       for _, dt := range []arrow.DataType{arrow.BinaryTypes.String, 
arrow.PrimitiveTypes.Int32,
+               arrow.Null, arrow.FixedWidthTypes.Boolean} {
+               _, err := countFn.DispatchBest(dt)
+               assert.NoErrorf(t, err, "count should accept %s", dt)
+       }
+}
+
+func TestScalarAggregateFunctionAddKernel(t *testing.T) {
+       fn := compute.NewScalarAggregateFunction("test_add_kernel", 
compute.Unary(), compute.EmptyFuncDoc)
+
+       in := []exec.InputType{exec.NewExactInput(arrow.PrimitiveTypes.Int64)}
+       out := exec.NewOutputType(arrow.PrimitiveTypes.Int64)
+       init := func(*exec.KernelCtx, exec.KernelInitArgs) (exec.KernelState, 
error) { return struct{}{}, nil }
+       consume := func(*exec.KernelCtx, *exec.ExecSpan) error { return nil }
+       merge := func(*exec.KernelCtx, exec.KernelState, exec.KernelState) 
error { return nil }
+       finalize := func(*exec.KernelCtx) (*exec.AggregateResult, error) { 
return nil, nil }
+
+       // the arity has to match
+       assert.ErrorIs(t, fn.AddNewKernel(nil, out, init, consume, merge, 
finalize), arrow.ErrInvalid)
+
+       // every one of the four lifecycle functions is required
+       assert.ErrorIs(t, fn.AddNewKernel(in, out, nil, consume, merge, 
finalize), arrow.ErrInvalid)
+       assert.ErrorIs(t, fn.AddNewKernel(in, out, init, nil, merge, finalize), 
arrow.ErrInvalid)
+       assert.ErrorIs(t, fn.AddNewKernel(in, out, init, consume, nil, 
finalize), arrow.ErrInvalid)
+       assert.ErrorIs(t, fn.AddNewKernel(in, out, init, consume, merge, nil), 
arrow.ErrInvalid)
+       assert.Zero(t, fn.NumKernels())
+
+       require.NoError(t, fn.AddNewKernel(in, out, init, consume, merge, 
finalize))
+       assert.Equal(t, 1, fn.NumKernels())
+       assert.False(t, fn.Kernels()[0].IsOrdered())
+}
+
+func TestScalarAggregateOptionsForms(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+
+       const vals = `[1, null, 3, -2]`
+       // nil options select the defaults: skip nulls, min_count 1
+       assertInt64Result(t, aggOf(t, ctx, "sum", nil, mem, 
arrow.PrimitiveTypes.Int64, vals), 2)
+       assertInt64Result(t, aggOf(t, ctx, "sum", 
compute.DefaultScalarAggregateOptions(), mem, arrow.PrimitiveTypes.Int64, 
vals), 2)
+       assert.Equal(t, &compute.ScalarAggregateOptions{SkipNulls: true, 
MinCount: 1}, compute.DefaultScalarAggregateOptions())
+       assert.Equal(t, &compute.CountOptions{Mode: compute.CountOnlyValid}, 
compute.DefaultCountOptions())
+
+       // an explicitly supplied zero value stays {false, 0} and is not
+       // silently turned into the defaults
+       assertNullResult(t, aggOf(t, ctx, "sum", 
&compute.ScalarAggregateOptions{}, mem, arrow.PrimitiveTypes.Int64, vals),
+               arrow.PrimitiveTypes.Int64)
+
+       // options are accepted by pointer and by value
+       arr := aggArray(t, mem, arrow.PrimitiveTypes.Int64, vals)
+       defer arr.Release()
+       d := compute.NewDatumWithoutOwning(arr)
+       assertInt64Result(t, callAgg(t, ctx, "sum", *skipNulls(true, 1), d), 2)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
compute.CountOptions{Mode: compute.CountAllRows}, d), 4)
+
+       // and the option types are registered for deserialization
+       assert.Equal(t, "ScalarAggregateOptions", 
compute.ScalarAggregateOptions{}.TypeName())
+       assert.Equal(t, "CountOptions", compute.CountOptions{}.TypeName())
+
+       _, err := compute.CallFunction(ctx, "sum", &compute.CountOptions{}, d)
+       assert.ErrorIs(t, err, arrow.ErrInvalid)
+       _, err = compute.CallFunction(ctx, "count", &compute.CountOptions{Mode: 
compute.CountMode(17)}, d)
+       assert.ErrorIs(t, err, arrow.ErrInvalid)
+}
+
+func TestScalarAggregateSumNumericTypes(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+
+       for _, dt := range []arrow.DataType{arrow.PrimitiveTypes.Int8, 
arrow.PrimitiveTypes.Int16,
+               arrow.PrimitiveTypes.Int32, arrow.PrimitiveTypes.Int64} {
+               t.Run(dt.String(), func(t *testing.T) {
+                       sc := aggOf(t, ctx, "sum", nil, mem, dt, `[1, null, 3, 
-2]`)
+                       assertInt64Result(t, sc, 2)
+                       assertInt64Result(t, aggOf(t, ctx, "count", nil, mem, 
dt, `[1, null, 3, -2]`), 3)
+               })
+       }
+
+       for _, dt := range []arrow.DataType{arrow.PrimitiveTypes.Uint8, 
arrow.PrimitiveTypes.Uint16,
+               arrow.PrimitiveTypes.Uint32, arrow.PrimitiveTypes.Uint64} {
+               t.Run(dt.String(), func(t *testing.T) {
+                       assertUint64Result(t, aggOf(t, ctx, "sum", nil, mem, 
dt, `[1, null, 3, 7]`), 11)
+                       assertInt64Result(t, aggOf(t, ctx, "count", nil, mem, 
dt, `[1, null, 3, 7]`), 3)
+               })
+       }
+
+       for _, dt := range []arrow.DataType{arrow.PrimitiveTypes.Float32, 
arrow.PrimitiveTypes.Float64} {
+               t.Run(dt.String(), func(t *testing.T) {
+                       assertFloat64Result(t, aggOf(t, ctx, "sum", nil, mem, 
dt, `[1.5, null, -2.5, 4.0]`), 3.0)
+                       assertInt64Result(t, aggOf(t, ctx, "count", nil, mem, 
dt, `[1.5, null, -2.5, 4.0]`), 3)
+               })
+       }
+
+       t.Run("bool", func(t *testing.T) {
+               const vals = `[true, null, false, true]`
+               assertUint64Result(t, aggOf(t, ctx, "sum", nil, mem, 
arrow.FixedWidthTypes.Boolean, vals), 2)
+               assertInt64Result(t, aggOf(t, ctx, "count", nil, mem, 
arrow.FixedWidthTypes.Boolean, vals), 3)
+       })
+}
+
+func TestScalarAggregateSumNullType(t *testing.T) {
+       ctx, _ := aggTestContext(t)
+
+       nulls := array.NewNull(3)
+       defer nulls.Release()
+       d := compute.NewDatumWithoutOwning(nulls)
+
+       assertNullResult(t, callAgg(t, ctx, "sum", nil, d), 
arrow.PrimitiveTypes.Int64)
+       assertInt64Result(t, callAgg(t, ctx, "sum", skipNulls(true, 0), d), 0)
+       assertNullResult(t, callAgg(t, ctx, "sum", skipNulls(false, 0), d), 
arrow.PrimitiveTypes.Int64)
+
+       assertInt64Result(t, callAgg(t, ctx, "count", nil, d), 0)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountOnlyNull}, d), 3)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountAllRows}, d), 3)
+
+       // an empty null array never saw a value, so even skip_nulls=false sums
+       // to zero when min_count allows it
+       empty := array.NewNull(0)
+       defer empty.Release()
+       assertInt64Result(t, callAgg(t, ctx, "sum", skipNulls(false, 0), 
compute.NewDatumWithoutOwning(empty)), 0)
+}
+
+func TestScalarAggregateSumNoNullsAllNullEmpty(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+       i64 := arrow.PrimitiveTypes.Int64
+
+       t.Run("no nulls", func(t *testing.T) {
+               const vals = `[4, 5, 6]`
+               assertInt64Result(t, aggOf(t, ctx, "sum", nil, mem, i64, vals), 
15)
+               assertInt64Result(t, aggOf(t, ctx, "sum", skipNulls(true, 0), 
mem, i64, vals), 15)
+               assertInt64Result(t, aggOf(t, ctx, "sum", skipNulls(false, 1), 
mem, i64, vals), 15)
+               assertInt64Result(t, aggOf(t, ctx, "sum", skipNulls(false, 0), 
mem, i64, vals), 15)
+               assertInt64Result(t, aggOf(t, ctx, "count", nil, mem, i64, 
vals), 3)
+               assertInt64Result(t, aggOf(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountOnlyNull}, mem, i64, vals), 0)
+               assertInt64Result(t, aggOf(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountAllRows}, mem, i64, vals), 3)
+       })
+
+       t.Run("all null", func(t *testing.T) {
+               const vals = `[null, null]`
+               assertNullResult(t, aggOf(t, ctx, "sum", nil, mem, i64, vals), 
i64)
+               assertInt64Result(t, aggOf(t, ctx, "sum", skipNulls(true, 0), 
mem, i64, vals), 0)
+               assertNullResult(t, aggOf(t, ctx, "sum", skipNulls(false, 1), 
mem, i64, vals), i64)
+               assertNullResult(t, aggOf(t, ctx, "sum", skipNulls(false, 0), 
mem, i64, vals), i64)
+               assertInt64Result(t, aggOf(t, ctx, "count", nil, mem, i64, 
vals), 0)
+               assertInt64Result(t, aggOf(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountOnlyNull}, mem, i64, vals), 2)
+               assertInt64Result(t, aggOf(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountAllRows}, mem, i64, vals), 2)
+       })
+
+       t.Run("empty", func(t *testing.T) {
+               const vals = `[]`
+               assertNullResult(t, aggOf(t, ctx, "sum", nil, mem, i64, vals), 
i64)
+               assertInt64Result(t, aggOf(t, ctx, "sum", skipNulls(true, 0), 
mem, i64, vals), 0)
+               assertNullResult(t, aggOf(t, ctx, "sum", skipNulls(false, 1), 
mem, i64, vals), i64)
+               // nothing was consumed, so no null was observed either
+               assertInt64Result(t, aggOf(t, ctx, "sum", skipNulls(false, 0), 
mem, i64, vals), 0)
+               assertInt64Result(t, aggOf(t, ctx, "count", nil, mem, i64, 
vals), 0)
+               assertInt64Result(t, aggOf(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountOnlyNull}, mem, i64, vals), 0)
+               assertInt64Result(t, aggOf(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountAllRows}, mem, i64, vals), 0)
+       })
+}
+
+func TestScalarAggregateSkipNullsAndMinCount(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+       i64 := arrow.PrimitiveTypes.Int64
+       const vals = `[1, null, 3, -2]`
+
+       assertNullResult(t, aggOf(t, ctx, "sum", skipNulls(false, 1), mem, i64, 
vals), i64)
+       assertInt64Result(t, aggOf(t, ctx, "sum", skipNulls(true, 3), mem, i64, 
vals), 2)
+       assertNullResult(t, aggOf(t, ctx, "sum", skipNulls(true, 4), mem, i64, 
vals), i64)
+
+       // count ignores ScalarAggregateOptions entirely
+       assertInt64Result(t, aggOf(t, ctx, "count", nil, mem, i64, vals), 3)
+       assertInt64Result(t, aggOf(t, ctx, "count", &compute.CountOptions{Mode: 
compute.CountOnlyNull}, mem, i64, vals), 1)
+       assertInt64Result(t, aggOf(t, ctx, "count", &compute.CountOptions{Mode: 
compute.CountAllRows}, mem, i64, vals), 4)
+}
+
+func TestScalarAggregateSumWraps(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+
+       assertInt64Result(t, aggOf(t, ctx, "sum", nil, mem, 
arrow.PrimitiveTypes.Int64,
+               fmt.Sprintf(`[%d, 1]`, int64(math.MaxInt64))), math.MinInt64)
+       assertUint64Result(t, aggOf(t, ctx, "sum", nil, mem, 
arrow.PrimitiveTypes.Uint64,
+               fmt.Sprintf(`[%d, 2]`, uint64(math.MaxUint64))), 1)
+       assertUint64Result(t, aggOf(t, ctx, "sum", nil, mem, 
arrow.PrimitiveTypes.Uint64,
+               fmt.Sprintf(`[%d, 1]`, uint64(math.MaxUint64))), 0)
+       // a narrow integer widens into its accumulator instead of wrapping
+       assertInt64Result(t, aggOf(t, ctx, "sum", nil, mem, 
arrow.PrimitiveTypes.Int8, `[127, 1]`), 128)
+}
+
+func TestScalarAggregateSumFloatNaN(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+
+       sc := aggOf(t, ctx, "sum", nil, mem, arrow.PrimitiveTypes.Float64, 
`[null, 1.0]`)
+       assertFloat64Result(t, sc, 1.0)
+
+       arr := aggArray(t, mem, arrow.PrimitiveTypes.Float64, `[1.0, 2.0]`)
+       defer arr.Release()
+       nanArr := array.NewFloat64Builder(mem)
+       defer nanArr.Release()
+       nanArr.AppendValues([]float64{math.NaN(), 1.0}, nil)
+       withNaN := nanArr.NewArray()
+       defer withNaN.Release()
+
+       res := callAgg(t, ctx, "sum", nil, 
compute.NewDatumWithoutOwning(withNaN))
+       v, ok := res.(*scalar.Float64)
+       require.True(t, ok)
+       assert.True(t, math.IsNaN(v.Value))
+}
+
+func TestScalarAggregateCountTypes(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+
+       strs := aggArray(t, mem, arrow.BinaryTypes.String, `["a", null, "c"]`)
+       defer strs.Release()
+       d := compute.NewDatumWithoutOwning(strs)
+       assertInt64Result(t, callAgg(t, ctx, "count", nil, d), 2)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountOnlyNull}, d), 1)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountAllRows}, d), 3)
+}
+
+// The logical null count of a run-end encoded or dictionary array is not the
+// number of unset bits in its own validity bitmap, so count rejects those
+// inputs rather than returning a wrong answer.
+func TestScalarAggregateCountLogicalNulls(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+
+       dictType := &arrow.DictionaryType{IndexType: arrow.PrimitiveTypes.Int8, 
ValueType: arrow.BinaryTypes.String}
+       bldr := array.NewDictionaryBuilder(mem, dictType)
+       defer bldr.Release()
+       require.NoError(t, bldr.AppendValueFromString("a"))
+       bldr.AppendNull()
+       dict := bldr.NewArray()
+       defer dict.Release()
+
+       d := compute.NewDatumWithoutOwning(dict)
+       _, err := compute.CallFunction(ctx, "count", nil, d)
+       assert.ErrorIs(t, err, arrow.ErrNotImplemented)
+
+       // counting every row does not need a null count at all
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountAllRows}, d), 2)
+}
+
+func TestScalarAggregateScalarInput(t *testing.T) {
+       ctx, _ := aggTestContext(t)
+
+       five := compute.NewDatum(scalar.NewInt64Scalar(5))
+       defer five.Release()
+       assertInt64Result(t, callAgg(t, ctx, "sum", nil, five), 5)
+       assertInt64Result(t, callAgg(t, ctx, "count", nil, five), 1)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountOnlyNull}, five), 0)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountAllRows}, five), 1)
+
+       nullScalar := 
compute.NewDatum(scalar.MakeNullScalar(arrow.PrimitiveTypes.Int64))
+       defer nullScalar.Release()
+       assertNullResult(t, callAgg(t, ctx, "sum", nil, nullScalar), 
arrow.PrimitiveTypes.Int64)
+       assertInt64Result(t, callAgg(t, ctx, "sum", skipNulls(true, 0), 
nullScalar), 0)
+       assertInt64Result(t, callAgg(t, ctx, "count", nil, nullScalar), 0)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountOnlyNull}, nullScalar), 1)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountAllRows}, nullScalar), 1)
+
+       boolTrue := compute.NewDatum(scalar.NewBooleanScalar(true))
+       defer boolTrue.Release()
+       assertUint64Result(t, callAgg(t, ctx, "sum", nil, boolTrue), 1)
+
+       // a narrow input is still summed into its accumulator type
+       three := compute.NewDatum(scalar.NewInt8Scalar(3))
+       defer three.Release()
+       assertInt64Result(t, callAgg(t, ctx, "sum", nil, three), 3)
+       assertInt64Result(t, callAgg(t, ctx, "count", nil, three), 1)
+}
+
+func TestScalarAggregateChunked(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+       i64 := arrow.PrimitiveTypes.Int64
+
+       chunks := make([]arrow.Array, 0, 3)
+       for _, vals := range []string{`[1, 2]`, `[null, 4]`, `[5]`} {
+               chunks = append(chunks, aggArray(t, mem, i64, vals))
+       }
+       chunked := arrow.NewChunked(i64, chunks)
+       for _, c := range chunks {
+               c.Release()
+       }
+       defer chunked.Release()
+
+       d := compute.NewDatumWithoutOwning(chunked)
+       assertInt64Result(t, callAgg(t, ctx, "sum", nil, d), 12)
+       assertNullResult(t, callAgg(t, ctx, "sum", skipNulls(false, 1), d), i64)
+       assertInt64Result(t, callAgg(t, ctx, "count", nil, d), 4)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountOnlyNull}, d), 1)
+       assertInt64Result(t, callAgg(t, ctx, "count", 
&compute.CountOptions{Mode: compute.CountAllRows}, d), 5)
+
+       emptyChunked := arrow.NewChunked(i64, nil)
+       defer emptyChunked.Release()
+       ed := compute.NewDatumWithoutOwning(emptyChunked)
+       assertNullResult(t, callAgg(t, ctx, "sum", nil, ed), i64)
+       assertInt64Result(t, callAgg(t, ctx, "sum", skipNulls(true, 0), ed), 0)
+       assertInt64Result(t, callAgg(t, ctx, "count", nil, ed), 0)
+}
+
+// xorshiftDoubles reproduces the fixture the C++ results were taken from.
+func xorshiftDoubles(n int) []float64 {
+       s := uint64(0x9E3779B97F4A7C15)
+       out := make([]float64, n)
+       for i := range out {
+               s ^= s >> 12
+               s ^= s << 25
+               s ^= s >> 27
+               r := s * 0x2545F4914F6CDD1D
+               u := float64(r>>11) / float64(uint64(1)<<53)
+               out[i] = (u - 0.5) * 1000.0
+       }
+       return out
+}
+
+// TestScalarAggregatePairwiseSum pins the float sum to the pairwise summation
+// the C++ implementation uses. A naive left to right sum of the same fixture
+// gives 75247.35643756675, so this test distinguishes the two.
+func TestScalarAggregatePairwiseSum(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+
+       vals := xorshiftDoubles(100001)
+       valid := make([]bool, len(vals))
+       for i := range valid {
+               valid[i] = i%7 != 6
+       }
+
+       f64Bldr := array.NewFloat64Builder(mem)
+       defer f64Bldr.Release()
+       f64Bldr.AppendValues(vals, nil)
+       f64 := f64Bldr.NewArray()
+       defer f64.Release()
+
+       f64Bldr.AppendValues(vals, valid)
+       f64n := f64Bldr.NewArray()
+       defer f64n.Release()
+
+       f32Bldr := array.NewFloat32Builder(mem)
+       defer f32Bldr.Release()
+       f32vals := make([]float32, len(vals))
+       i32vals := make([]int32, len(vals))
+       for i, v := range vals {
+               f32vals[i] = float32(v)
+               i32vals[i] = int32(v)
+       }
+       f32Bldr.AppendValues(f32vals, nil)
+       f32 := f32Bldr.NewArray()
+       defer f32.Release()
+
+       i32Bldr := array.NewInt32Builder(mem)
+       defer i32Bldr.Release()
+       i32Bldr.AppendValues(i32vals, nil)
+       i32 := i32Bldr.NewArray()
+       defer i32.Release()
+
+       assertFloat64Result(t, callAgg(t, ctx, "sum", nil, 
compute.NewDatumWithoutOwning(f64)), 75247.35643756694)
+
+       // the pinned value is the pairwise one: a naive left-to-right sum of 
the
+       // same fixture lands on a different double, so the assertions above can
+       // only pass with the C++ summation order
+       var naive float64
+       for _, v := range vals {
+               naive += v
+       }
+       assert.NotEqual(t, 75247.35643756694, naive)
+       assertFloat64Result(t, callAgg(t, ctx, "sum", nil, 
compute.NewDatumWithoutOwning(f64n)), 94384.73372326966)
+       assertFloat64Result(t, callAgg(t, ctx, "sum", nil, 
compute.NewDatumWithoutOwning(f32)), 75247.3534175111)
+       assertInt64Result(t, callAgg(t, ctx, "sum", nil, 
compute.NewDatumWithoutOwning(i32)), 75006)
+
+       // slicing shifts the offset without changing the summation order
+       sliced := array.NewSlice(f64, 5, int64(f64.Len()))
+       defer sliced.Release()
+       assertFloat64Result(t, callAgg(t, ctx, "sum", nil, 
compute.NewDatumWithoutOwning(sliced)), 75650.35881616312)
+
+       slicedNulls := array.NewSlice(f64n, 5, int64(f64n.Len()))
+       defer slicedNulls.Release()
+       assertFloat64Result(t, callAgg(t, ctx, "sum", nil, 
compute.NewDatumWithoutOwning(slicedNulls)), 94787.73610186583)
+}
+
+// TestScalarAggregateChunkSizeInvariance runs every aggregate over the same
+// input at chunk sizes 1 to 65 and compares against the single-span result.
+// A kernel which cached a null count back into the span the executor re-slices
+// would fail here without failing any of the tests above.
+func TestScalarAggregateChunkSizeInvariance(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+
+       cases := []struct {
+               dt   arrow.DataType
+               vals string
+       }{
+               {arrow.PrimitiveTypes.Int64, `[1, null, 3, -2, 5, null, 7, 8, 
null, 10, 11, 12, null, 14, 15, 16, 17, null, 19, 20]`},
+               {arrow.PrimitiveTypes.Uint32, `[1, null, 3, 7, 5, null, 7, 8, 
null, 10, 11, 12, null, 14, 15, 16, 17, null, 19, 20]`},
+               {arrow.PrimitiveTypes.Float64, `[1.5, null, -2.5, 4.0, 5.5, 
null, 7.25, 8.125, null, 10.0, 11, 12, null, 14, 15, 16, 17, null, 19, 20]`},
+               {arrow.FixedWidthTypes.Boolean, `[true, null, false, true, 
true, null, false, true, null, true, false, true, null, true, true, false, 
true, null, false, true]`},
+               {arrow.BinaryTypes.String, `["a", null, "c", "d", "e", null, 
"g", "h", null, "j", "k", "l", null, "n", "o", "p", "q", null, "s", "t"]`},
+       }
+
+       optSets := []compute.FunctionOptions{nil, skipNulls(true, 0), 
skipNulls(false, 1), skipNulls(true, 21)}
+       countOpts := []compute.FunctionOptions{
+               &compute.CountOptions{Mode: compute.CountOnlyValid},
+               &compute.CountOptions{Mode: compute.CountOnlyNull},
+               &compute.CountOptions{Mode: compute.CountAllRows},
+       }
+
+       for _, tc := range cases {
+               arr := aggArray(t, mem, tc.dt, tc.vals)
+               d := compute.NewDatumWithoutOwning(arr)
+
+               type invocation struct {
+                       fn   string
+                       opts compute.FunctionOptions
+               }
+               invocations := make([]invocation, 0, 
len(optSets)+len(countOpts))
+               for _, o := range countOpts {
+                       invocations = append(invocations, invocation{"count", 
o})
+               }
+               if tc.dt.ID() != arrow.STRING {
+                       for _, o := range optSets {
+                               invocations = append(invocations, 
invocation{"sum", o})
+                       }
+               }
+
+               for _, inv := range invocations {
+                       want := callAgg(t, ctx, inv.fn, inv.opts, d)
+                       for chunkSize := int64(1); chunkSize <= 65; chunkSize++ 
{
+                               ectx := compute.DefaultExecCtx()
+                               ectx.ChunkSize = chunkSize
+                               got := callAgg(t, compute.SetExecCtx(ctx, 
ectx), inv.fn, inv.opts, d)
+                               assert.Truef(t, scalar.Equals(want, got),
+                                       "%s over %s at chunk size %d: expected 
%s, got %s", inv.fn, tc.dt, chunkSize, want, got)
+                       }
+               }
+               arr.Release()
+       }
+}
+
+// TestScalarAggregateMerge consumes the two halves of an input into separate
+// states and merges them, checking that the merged state finalizes to the
+// same value the executor produces for the whole input.
+func TestScalarAggregateMerge(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+
+       cases := []struct {
+               fn   string
+               dt   arrow.DataType
+               vals string
+               opts any
+       }{
+               {"sum", arrow.PrimitiveTypes.Int64, `[1, null, 3, -2, 5, 6]`, 
skipNulls(true, 1)},
+               {"sum", arrow.PrimitiveTypes.Int64, `[1, null, 3, -2, 5, 6]`, 
skipNulls(false, 1)},
+               {"sum", arrow.PrimitiveTypes.Uint16, `[1, null, 3, 7, 5, 6]`, 
skipNulls(true, 1)},
+               {"sum", arrow.PrimitiveTypes.Float64, `[1.5, null, -2.5, 4.0, 
5.5, 6.25]`, skipNulls(true, 1)},
+               {"sum", arrow.FixedWidthTypes.Boolean, `[true, null, false, 
true, true, false]`, skipNulls(true, 1)},
+               {"count", arrow.PrimitiveTypes.Int64, `[1, null, 3, -2, 5, 6]`, 
&compute.CountOptions{Mode: compute.CountOnlyValid}},
+               {"count", arrow.PrimitiveTypes.Int64, `[1, null, 3, -2, 5, 6]`, 
&compute.CountOptions{Mode: compute.CountOnlyNull}},
+               {"count", arrow.PrimitiveTypes.Int64, `[1, null, 3, -2, 5, 6]`, 
&compute.CountOptions{Mode: compute.CountAllRows}},
+       }
+
+       reg := compute.GetFunctionRegistry()
+       for _, tc := range cases {
+               t.Run(tc.fn+"/"+tc.dt.String(), func(t *testing.T) {
+                       arr := aggArray(t, mem, tc.dt, tc.vals)
+                       defer arr.Release()
+
+                       fn, ok := reg.GetFunction(tc.fn)
+                       require.True(t, ok)
+                       k, err := fn.DispatchBest(tc.dt)
+                       require.NoError(t, err)
+                       aggKernel := k.(*exec.ScalarAggKernel)
+
+                       kctx := &exec.KernelCtx{Ctx: ctx, Kernel: aggKernel}
+                       initArgs := exec.KernelInitArgs{Kernel: aggKernel, 
Inputs: []arrow.DataType{tc.dt}, Options: tc.opts}
+
+                       states := make([]exec.KernelState, 0, 2)
+                       for _, half := range [][2]int64{{0, 3}, {3, 6}} {
+                               state, err := aggKernel.GetInitFn()(kctx, 
initArgs)
+                               require.NoError(t, err)
+
+                               slice := array.NewSlice(arr, half[0], half[1])
+                               span := exec.ExecSpan{Len: int64(slice.Len()), 
Values: make([]exec.ExecValue, 1)}
+                               span.Values[0].Array.SetMembers(slice.Data())
+
+                               stateCtx := &exec.KernelCtx{Ctx: ctx, Kernel: 
aggKernel, State: state}
+                               require.NoError(t, aggKernel.Consume(stateCtx, 
&span))
+                               slice.Release()
+                               states = append(states, state)
+                       }
+
+                       merged, err := exec.MergeAll(aggKernel, kctx, states)
+                       require.NoError(t, err)
+                       kctx.State = merged
+                       res, err := aggKernel.Finalize(kctx)
+                       require.NoError(t, err)
+                       defer res.Release()
+
+                       var opts compute.FunctionOptions
+                       if o, ok := tc.opts.(compute.FunctionOptions); ok {
+                               opts = o
+                       }
+                       want := callAgg(t, ctx, tc.fn, opts, 
compute.NewDatumWithoutOwning(arr))
+                       assert.Truef(t, scalar.Equals(want, res.Scalar()), 
"expected %s, got %s", want, res.Scalar())
+               })
+       }
+}
+
+func TestScalarAggregateWrappers(t *testing.T) {
+       ctx, mem := aggTestContext(t)
+
+       arr := aggArray(t, mem, arrow.PrimitiveTypes.Int64, `[1, null, 3, -2]`)
+       defer arr.Release()
+       d := compute.NewDatumWithoutOwning(arr)
+
+       res, err := compute.Sum(ctx, *compute.DefaultScalarAggregateOptions(), 
d)
+       require.NoError(t, err)
+       assertInt64Result(t, res.(*compute.ScalarDatum).Value, 2)
+       res.Release()
+
+       res, err = compute.Count(ctx, *compute.DefaultCountOptions(), d)
+       require.NoError(t, err)
+       assertInt64Result(t, res.(*compute.ScalarDatum).Value, 3)
+       res.Release()
+
+       res, err = compute.Count(ctx, compute.CountOptions{Mode: 
compute.CountAllRows}, d)
+       require.NoError(t, err)
+       assertInt64Result(t, res.(*compute.ScalarDatum).Value, 4)
+       res.Release()
+}
+
+// ----------------------------------------------------------------------
+// framework tests which the primitive count and sum kernels cannot expose
+
+var errTestAgg = errors.New("test aggregate failure")
+
+type testAggBehaviour struct {
+       consumeErr  error
+       finalizeErr error
+       initErr     error
+       cleanupErr  error
+       arrayResult bool
+       cleanups    *int
+}
+
+type testAggState struct {
+       behaviour *testAggBehaviour
+       mem       memory.Allocator
+       // buf is owned by the state and released by the cleanup function; the
+       // result finalize returns must stay valid after that
+       buf   *memory.Buffer
+       total int64
+}
+
+// newTestAggFunction registers an aggregate function whose result is backed
+// by allocated memory, so that the ownership rule of Finalize is exercised.
+func newTestAggFunction(t *testing.T, name string, behaviour 
*testAggBehaviour, mem memory.Allocator) *compute.ScalarAggregateFunction {
+       t.Helper()
+
+       outType := arrow.DataType(arrow.BinaryTypes.String)
+       if behaviour.arrayResult {
+               outType = arrow.PrimitiveTypes.Int64
+       }
+
+       fn := compute.NewScalarAggregateFunction(name, compute.Unary(), 
compute.EmptyFuncDoc)
+       kernel := exec.NewScalarAggKernel(
+               
[]exec.InputType{exec.NewExactInput(arrow.PrimitiveTypes.Int64)},
+               exec.NewOutputType(outType),
+               func(*exec.KernelCtx, exec.KernelInitArgs) (exec.KernelState, 
error) {
+                       if behaviour.initErr != nil {
+                               return nil, behaviour.initErr
+                       }
+                       buf := memory.NewResizableBuffer(mem)
+                       buf.Resize(8)
+                       return &testAggState{behaviour: behaviour, mem: mem, 
buf: buf}, nil
+               },
+               func(ctx *exec.KernelCtx, span *exec.ExecSpan) error {
+                       st := ctx.State.(*testAggState)
+                       if st.behaviour.consumeErr != nil {
+                               return st.behaviour.consumeErr
+                       }
+                       st.total += span.Len
+                       return nil
+               },
+               func(_ *exec.KernelCtx, src, dst exec.KernelState) error {
+                       dst.(*testAggState).total += src.(*testAggState).total
+                       return nil
+               },
+               func(ctx *exec.KernelCtx) (*exec.AggregateResult, error) {
+                       st := ctx.State.(*testAggState)
+                       if st.behaviour.finalizeErr != nil {
+                               return nil, st.behaviour.finalizeErr
+                       }
+                       if st.behaviour.arrayResult {
+                               bldr := array.NewInt64Builder(st.mem)
+                               defer bldr.Release()
+                               bldr.Append(st.total)
+                               arr := bldr.NewArray()
+                               defer arr.Release()
+                               arr.Data().Retain()
+                               return exec.NewArrayResult(arr.Data()), nil
+                       }
+                       // a buffer-backed scalar which must outlive the 
state's own
+                       // buffer: the result takes its own reference
+                       out := memory.NewResizableBuffer(st.mem)
+                       defer out.Release()
+                       text := fmt.Sprintf("total=%d", st.total)
+                       out.Resize(len(text))
+                       copy(out.Bytes(), text)
+                       return 
exec.NewScalarResult(scalar.NewStringScalarFromBuffer(out)), nil
+               })
+       kernel.CleanupFn = func(_ *exec.KernelCtx, state exec.KernelState) 
error {
+               if state != nil {
+                       state.(*testAggState).buf.Release()
+               }
+               *behaviour.cleanups++
+               return behaviour.cleanupErr
+       }
+       require.NoError(t, fn.AddKernel(kernel))
+       return fn
+}
+
+func testAggContext(t *testing.T, ctx context.Context, fn compute.Function) 
context.Context {
+       t.Helper()
+       reg := compute.NewChildRegistry(compute.GetFunctionRegistry())
+       require.True(t, reg.AddFunction(fn, false))
+       ectx := compute.DefaultExecCtx()
+       ectx.Registry = reg
+       return compute.SetExecCtx(ctx, ectx)
+}
+
+// TestScalarAggregateResultOwnership checks that the value finalize returned
+// outlives the cleanup of the aggregate state which produced it, both for a
+// buffer-backed scalar and for an array result.
+func TestScalarAggregateResultOwnership(t *testing.T) {
+       for _, arrayResult := range []bool{false, true} {
+               name := "scalar result"
+               if arrayResult {
+                       name = "array result"
+               }
+               t.Run(name, func(t *testing.T) {
+                       ctx, mem := aggTestContext(t)
+
+                       var cleanups int
+                       behaviour := &testAggBehaviour{arrayResult: 
arrayResult, cleanups: &cleanups}
+                       fn := newTestAggFunction(t, "test_agg_ownership", 
behaviour, mem)
+                       ctx = testAggContext(t, ctx, fn)
+
+                       arr := aggArray(t, mem, arrow.PrimitiveTypes.Int64, 
`[1, 2, 3, 4, 5]`)
+                       defer arr.Release()
+
+                       res, err := compute.CallFunction(ctx, 
"test_agg_ownership", nil, compute.NewDatumWithoutOwning(arr))
+                       require.NoError(t, err)
+                       assert.Equal(t, 1, cleanups, "the aggregate state is 
cleaned up exactly once")
+
+                       // the state's buffer is gone by now; the result is not
+                       if arrayResult {
+                               ad := res.(*compute.ArrayDatum)
+                               out := ad.MakeArray()
+                               assert.Equal(t, []int64{5}, 
out.(*array.Int64).Int64Values())
+                               out.Release()
+                       } else {
+                               assert.Equal(t, "total=5", 
res.(*compute.ScalarDatum).Value.(*scalar.String).String())
+                       }
+                       res.Release()
+               })
+       }
+}
+
+// TestScalarAggregateStateCleanup checks the promise that cleanup runs
+// exactly once whatever happens.
+func TestScalarAggregateStateCleanup(t *testing.T) {
+       newCall := func(t *testing.T, behaviour *testAggBehaviour, vals string) 
(context.Context, compute.Datum, func()) {
+               ctx, mem := aggTestContext(t)
+               fn := newTestAggFunction(t, "test_agg_cleanup", behaviour, mem)
+               ctx = testAggContext(t, ctx, fn)
+               arr := aggArray(t, mem, arrow.PrimitiveTypes.Int64, vals)
+               return ctx, compute.NewDatumWithoutOwning(arr), arr.Release
+       }
+
+       t.Run("success", func(t *testing.T) {
+               var cleanups int
+               ctx, d, done := newCall(t, &testAggBehaviour{cleanups: 
&cleanups}, `[1, 2, 3]`)
+               defer done()
+               res, err := compute.CallFunction(ctx, "test_agg_cleanup", nil, 
d)
+               require.NoError(t, err)
+               res.Release()
+               assert.Equal(t, 1, cleanups)
+       })
+
+       t.Run("empty input", func(t *testing.T) {
+               var cleanups int
+               ctx, d, done := newCall(t, &testAggBehaviour{cleanups: 
&cleanups}, `[]`)
+               defer done()
+               res, err := compute.CallFunction(ctx, "test_agg_cleanup", nil, 
d)
+               require.NoError(t, err)
+               res.Release()
+               assert.Equal(t, 1, cleanups)
+       })
+
+       t.Run("consume error", func(t *testing.T) {
+               var cleanups int
+               ctx, d, done := newCall(t, &testAggBehaviour{consumeErr: 
errTestAgg, cleanups: &cleanups}, `[1, 2, 3]`)
+               defer done()
+               _, err := compute.CallFunction(ctx, "test_agg_cleanup", nil, d)
+               assert.ErrorIs(t, err, errTestAgg)
+               assert.Equal(t, 1, cleanups)
+       })
+
+       t.Run("finalize error", func(t *testing.T) {
+               var cleanups int
+               ctx, d, done := newCall(t, &testAggBehaviour{finalizeErr: 
errTestAgg, cleanups: &cleanups}, `[1, 2, 3]`)
+               defer done()
+               _, err := compute.CallFunction(ctx, "test_agg_cleanup", nil, d)
+               assert.ErrorIs(t, err, errTestAgg)
+               assert.Equal(t, 1, cleanups)
+       })
+
+       t.Run("cleanup error is reported", func(t *testing.T) {
+               var cleanups int
+               ctx, d, done := newCall(t, &testAggBehaviour{cleanupErr: 
errTestAgg, cleanups: &cleanups}, `[1, 2, 3]`)
+               defer done()
+               _, err := compute.CallFunction(ctx, "test_agg_cleanup", nil, d)
+               assert.ErrorIs(t, err, errTestAgg)
+               assert.Equal(t, 1, cleanups)
+       })
+
+       t.Run("init error", func(t *testing.T) {
+               var cleanups int
+               ctx, d, done := newCall(t, &testAggBehaviour{initErr: 
errTestAgg, cleanups: &cleanups}, `[1, 2, 3]`)
+               defer done()
+               _, err := compute.CallFunction(ctx, "test_agg_cleanup", nil, d)
+               assert.ErrorIs(t, err, errTestAgg)
+               // no state was produced, so there is nothing to clean up
+               assert.Zero(t, cleanups)
+       })
+
+       t.Run("cancellation", func(t *testing.T) {
+               var cleanups int
+               ctx, d, done := newCall(t, &testAggBehaviour{cleanups: 
&cleanups}, `[1, 2, 3]`)
+               defer done()
+               cancelled, cancel := context.WithCancel(ctx)
+               cancel()
+               _, err := compute.CallFunction(cancelled, "test_agg_cleanup", 
nil, d)
+               assert.ErrorIs(t, err, context.Canceled)
+               assert.Equal(t, 1, cleanups)
+       })
+}
+
+// TestScalarAggregateScalarSpanLength drives the executor with a batch of
+// scalars whose logical length is greater than one, which the public
+// CallFunction path cannot produce. A kernel which ignored the span length
+// for scalar input would pass every other test in this file.
+func TestScalarAggregateScalarSpanLength(t *testing.T) {
+       ctx, _ := aggTestContext(t)
+
+       reg := compute.GetFunctionRegistry()
+       const spanLen = 5
+
+       run := func(t *testing.T, fname string, opts compute.FunctionOptions, 
sc scalar.Scalar) scalar.Scalar {
+               t.Helper()
+               fn, ok := reg.GetFunction(fname)
+               require.True(t, ok)
+               k, err := fn.DispatchBest(sc.DataType())
+               require.NoError(t, err)
+
+               kctx := &exec.KernelCtx{Ctx: ctx, Kernel: k}
+               initArgs := exec.KernelInitArgs{Kernel: k, Inputs: 
[]arrow.DataType{sc.DataType()}, Options: opts}
+               kctx.State, err = k.GetInitFn()(kctx, initArgs)
+               require.NoError(t, err)
+
+               executor := compute.NewScalarAggExecutor()
+               defer executor.Clear()
+               require.NoError(t, executor.Init(kctx, initArgs))
+
+               batch := &compute.ExecBatch{Values: 
[]compute.Datum{compute.NewDatumWithoutOwning(sc)}, Len: spanLen}
+               ch := make(chan compute.Datum, 1)
+               go func() {
+                       defer close(ch)
+                       require.NoError(t, executor.Execute(ctx, batch, ch))
+               }()
+               out := executor.WrapResults(ctx, ch, false)
+               require.NotNil(t, out)
+               defer out.Release()
+               return out.(*compute.ScalarDatum).Value
+       }
+
+       // the scalar counts once per row of the span
+       assertInt64Result(t, run(t, "sum", nil, scalar.NewInt64Scalar(3)), 
3*spanLen)
+       assertUint64Result(t, run(t, "sum", nil, scalar.NewUint8Scalar(2)), 
2*spanLen)
+       assertFloat64Result(t, run(t, "sum", nil, 
scalar.NewFloat64Scalar(1.5)), 1.5*spanLen)
+       assertUint64Result(t, run(t, "sum", nil, 
scalar.NewBooleanScalar(true)), spanLen)
+       assertInt64Result(t, run(t, "count", nil, scalar.NewInt64Scalar(3)), 
spanLen)
+       assertInt64Result(t, run(t, "count", &compute.CountOptions{Mode: 
compute.CountAllRows}, scalar.NewInt64Scalar(3)), spanLen)
+       assertInt64Result(t, run(t, "count", &compute.CountOptions{Mode: 
compute.CountOnlyNull}, scalar.NewInt64Scalar(3)), 0)
+
+       nullSc := scalar.MakeNullScalar(arrow.PrimitiveTypes.Int64)
+       assertInt64Result(t, run(t, "count", &compute.CountOptions{Mode: 
compute.CountOnlyNull}, nullSc), spanLen)
+       assertInt64Result(t, run(t, "count", nil, nullSc), 0)
+       assertNullResult(t, run(t, "sum", nil, nullSc), 
arrow.PrimitiveTypes.Int64)
+
+       // the same result whatever the chunk size the span is cut into
+       small := compute.DefaultExecCtx()
+       small.ChunkSize = 2
+       prevCtx := ctx
+       ctx = compute.SetExecCtx(prevCtx, small)
+       assertInt64Result(t, run(t, "sum", nil, scalar.NewInt64Scalar(3)), 
3*spanLen)
+       assertInt64Result(t, run(t, "count", nil, scalar.NewInt64Scalar(3)), 
spanLen)
+       ctx = prevCtx
+}

Reply via email to