Coverage Report

Created: 2026-10-01 17:01

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exprs/aggregate/aggregate_function_foreach.h
Line
Count
Source
1
// Licensed to the Apache Software Foundation (ASF) under one
2
// or more contributor license agreements.  See the NOTICE file
3
// distributed with this work for additional information
4
// regarding copyright ownership.  The ASF licenses this file
5
// to you under the Apache License, Version 2.0 (the
6
// "License"); you may not use this file except in compliance
7
// with the License.  You may obtain a copy of the License at
8
//
9
//   http://www.apache.org/licenses/LICENSE-2.0
10
//
11
// Unless required by applicable law or agreed to in writing,
12
// software distributed under the License is distributed on an
13
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14
// KIND, either express or implied.  See the License for the
15
// specific language governing permissions and limitations
16
// under the License.
17
// This file is copied from
18
// https://github.com/ClickHouse/ClickHouse/blob/master/src/AggregateFunctions/Combinators/AggregateFunctionForEach.h
19
// and modified by Doris
20
21
#pragma once
22
23
#include <utility>
24
25
#include "common/status.h"
26
#include "core/assert_cast.h"
27
#include "core/column/column_nullable.h"
28
#include "core/data_type/data_type_array.h"
29
#include "core/data_type/data_type_nullable.h"
30
#include "exec/common/arithmetic_overflow.h"
31
#include "exprs/aggregate/aggregate_function.h"
32
#include "exprs/function/array/function_array_utils.h"
33
34
namespace doris {
35
36
struct AggregateFunctionForEachData {
37
    size_t dynamic_array_size = 0;
38
    char* array_of_aggregate_datas = nullptr;
39
};
40
41
/** Adaptor for aggregate functions.
42
  * Adding -ForEach suffix to aggregate function
43
  *  will convert that aggregate function to a function, accepting arrays,
44
  *  and applies aggregation for each corresponding elements of arrays independently,
45
  *  returning arrays of aggregated values on corresponding positions.
46
  *
47
  * Example: sumForEach of:
48
  *  [1, 2],
49
  *  [3, 4, 5],
50
  *  [6, 7]
51
  * will return:
52
  *  [10, 13, 5]
53
  *
54
  * TODO Allow variable number of arguments.
55
  */
56
class AggregateFunctionForEach : public AggregateFunctionNonFinalBase,
57
                                 public IAggregateFunctionDataHelper<AggregateFunctionForEachData,
58
                                                                     AggregateFunctionForEach>,
59
                                 VarargsExpression,
60
                                 NullableAggregateFunction {
61
protected:
62
    using Base =
63
            IAggregateFunctionDataHelper<AggregateFunctionForEachData, AggregateFunctionForEach>;
64
65
    AggregateFunctionPtr nested_function;
66
    const size_t nested_size_of_data;
67
    const size_t num_arguments;
68
69
    AggregateFunctionForEachData& ensure_aggregate_data(AggregateDataPtr __restrict place,
70
174
                                                        size_t new_size, Arena& arena) const {
71
174
        AggregateFunctionForEachData& state = data(place);
72
73
        /// Ensure we have aggregate states for new_size elements, allocate
74
        /// from arena if needed. When reallocating, we can't copy the
75
        /// states to new buffer with memcpy, because they may contain pointers
76
        /// to themselves. In particular, this happens when a state contains
77
        /// a PODArrayWithStackMemory, which stores small number of elements
78
        /// inline. This is why we create new empty states in the new buffer,
79
        /// and merge the old states to them.
80
174
        size_t old_size = state.dynamic_array_size;
81
174
        if (old_size < new_size) {
82
172
            static constexpr size_t MAX_ARRAY_SIZE = 100 * 1000000000ULL;
83
172
            if (new_size > MAX_ARRAY_SIZE) {
84
0
                throw Exception(ErrorCode::INTERNAL_ERROR,
85
0
                                "Suspiciously large array size ({}) in -ForEach aggregate function",
86
0
                                new_size);
87
0
            }
88
89
172
            size_t allocation_size = 0;
90
172
            if (common::mul_overflow(new_size, nested_size_of_data, allocation_size)) {
91
0
                throw Exception(ErrorCode::INTERNAL_ERROR,
92
0
                                "Allocation size ({} * {}) overflows in -ForEach aggregate "
93
0
                                "function, but it should've been prevented by previous checks",
94
0
                                new_size, nested_size_of_data);
95
0
            }
96
97
172
            char* old_state = state.array_of_aggregate_datas;
98
99
172
            char* new_state =
100
172
                    arena.aligned_alloc(allocation_size, nested_function->align_of_data());
101
102
172
            size_t num_created = 0;
103
172
            try {
104
1.34k
                for (; num_created < new_size; ++num_created) {
105
1.16k
                    nested_function->create(&new_state[num_created * nested_size_of_data]);
106
1.16k
                }
107
108
175
                for (size_t i = 0; i < old_size; ++i) {
109
3
                    nested_function->merge(&new_state[i * nested_size_of_data],
110
3
                                           &old_state[i * nested_size_of_data], arena);
111
3
                }
112
172
            } catch (...) {
113
4
                for (size_t i = 0; i < num_created; ++i) {
114
3
                    nested_function->destroy(&new_state[i * nested_size_of_data]);
115
3
                }
116
117
1
                throw;
118
1
            }
119
120
172
            for (size_t i = 0; i < old_size; ++i) {
121
1
                nested_function->destroy(&old_state[i * nested_size_of_data]);
122
1
            }
123
124
171
            state.array_of_aggregate_datas = new_state;
125
171
            state.dynamic_array_size = new_size;
126
171
        }
127
128
173
        return state;
129
174
    }
130
131
private:
132
    struct PreparedBatchColumns {
133
        Columns owned_columns;
134
        std::vector<const IColumn*> nested_columns;
135
        const ColumnArray::Offsets64* offsets = nullptr;
136
    };
137
138
    void add_aligned(AggregateDataPtr __restrict place, const IColumn** nested,
139
16
                     const ColumnArray::Offsets64& offsets, size_t row_num, Arena& arena) const {
140
16
        const size_t begin = offsets[row_num - 1];
141
16
        const size_t end = offsets[row_num];
142
16
        AggregateFunctionForEachData& state = ensure_aggregate_data(place, end - begin, arena);
143
144
16
        char* nested_state = state.array_of_aggregate_datas;
145
68
        for (size_t i = begin; i < end; ++i) {
146
52
            nested_function->add(nested_state, nested, i, arena);
147
52
            nested_state += nested_size_of_data;
148
52
        }
149
16
    }
150
151
    // Validate all selected rows first. If their physical begins differ, compact every argument
152
    // once for the batch and give them shared offsets;
153
    // E.g. data [9,8,7,6]/[1,2,3,4,5], offsets [1,4]/[2,5], skip row 0
154
    // -> data [8,7,6]/[3,4,5], shared offsets [0,3].
155
    template <typename IsSelected>
156
    Columns normalize_batch_columns(size_t batch_size, const IColumn** columns,
157
17
                                    IsSelected&& is_selected) const {
158
17
        bool offsets_aligned = true;
159
        // First pass: reject size mismatches before allocating, and detect whether the
160
        // original nested columns can use one shared element index.
161
169
        for (size_t row = 0; row < batch_size; ++row) {
162
152
            if (!is_selected(row)) {
163
4
                continue;
164
4
            }
165
166
148
            const auto& first_array =
167
148
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0]);
168
148
            const auto& first_offsets = first_array.get_offsets();
169
148
            const size_t first_begin = first_offsets[row - 1];
170
148
            const size_t row_size = first_offsets[row] - first_begin;
171
284
            for (size_t i = 1; i < num_arguments; ++i) {
172
136
                const auto& array =
173
136
                        assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
174
136
                const auto& offsets = array.get_offsets();
175
136
                const size_t begin = offsets[row - 1];
176
136
                if (offsets[row] - begin != row_size) {
177
0
                    throw Exception(ErrorCode::INTERNAL_ERROR,
178
0
                                    "Arrays passed to {} aggregate function have different sizes",
179
0
                                    get_name());
180
0
                }
181
136
                offsets_aligned &= row_size == 0 || begin == first_begin;
182
136
            }
183
148
        }
184
185
        // Fast path: all visible rows already address the same physical nested positions.
186
17
        if (offsets_aligned) {
187
13
            return {};
188
13
        }
189
190
        // Slow path: build the shared offsets once, copy only selected row slices, and preserve
191
        // the outer row count by emitting repeated offsets for unselected rows.
192
4
        auto offsets_column = ColumnArray::ColumnOffsets::create();
193
4
        auto& normalized_offsets = offsets_column->get_data();
194
4
        normalized_offsets.reserve(batch_size);
195
4
        size_t normalized_offset = 0;
196
4
        const auto& first_offsets =
197
4
                assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0])
198
4
                        .get_offsets();
199
16
        for (size_t row = 0; row < batch_size; ++row) {
200
12
            if (is_selected(row)) {
201
8
                normalized_offset += first_offsets[row] - first_offsets[row - 1];
202
8
            }
203
12
            normalized_offsets.push_back(normalized_offset);
204
12
        }
205
4
        ColumnPtr shared_offsets = std::move(offsets_column);
206
207
4
        Columns normalized_columns(num_arguments);
208
12
        for (size_t i = 0; i < num_arguments; ++i) {
209
8
            const auto& array =
210
8
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
211
8
            auto nested = array.get_data().clone_empty();
212
32
            for (size_t row = 0; row < batch_size; ++row) {
213
24
                if (is_selected(row)) {
214
16
                    const auto& offsets = array.get_offsets();
215
16
                    const size_t begin = offsets[row - 1];
216
16
                    const size_t row_size = offsets[row] - begin;
217
16
                    nested->insert_range_from(array.get_data(), begin, row_size);
218
16
                }
219
24
            }
220
8
            ColumnPtr nested_column = std::move(nested);
221
8
            normalized_columns[i] = ColumnArray::create(nested_column, shared_offsets);
222
8
        }
223
4
        return normalized_columns;
224
17
    }
Unexecuted instantiation: _ZNK5doris24AggregateFunctionForEach23normalize_batch_columnsIZNKS0_9add_batchEmPPcmPPKNS_7IColumnERNS_5ArenaEbEUlmE_EESt6vectorINS_3COWIS4_E13immutable_ptrIS4_EESaISF_EEmS7_OT_
_ZNK5doris24AggregateFunctionForEach23normalize_batch_columnsIZNKS0_18add_batch_selectedEmPPcmPPKNS_7IColumnERNS_5ArenaEEUlmE_EESt6vectorINS_3COWIS4_E13immutable_ptrIS4_EESaISF_EEmS7_OT_
Line
Count
Source
157
2
                                    IsSelected&& is_selected) const {
158
2
        bool offsets_aligned = true;
159
        // First pass: reject size mismatches before allocating, and detect whether the
160
        // original nested columns can use one shared element index.
161
8
        for (size_t row = 0; row < batch_size; ++row) {
162
6
            if (!is_selected(row)) {
163
2
                continue;
164
2
            }
165
166
4
            const auto& first_array =
167
4
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0]);
168
4
            const auto& first_offsets = first_array.get_offsets();
169
4
            const size_t first_begin = first_offsets[row - 1];
170
4
            const size_t row_size = first_offsets[row] - first_begin;
171
8
            for (size_t i = 1; i < num_arguments; ++i) {
172
4
                const auto& array =
173
4
                        assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
174
4
                const auto& offsets = array.get_offsets();
175
4
                const size_t begin = offsets[row - 1];
176
4
                if (offsets[row] - begin != row_size) {
177
0
                    throw Exception(ErrorCode::INTERNAL_ERROR,
178
0
                                    "Arrays passed to {} aggregate function have different sizes",
179
0
                                    get_name());
180
0
                }
181
4
                offsets_aligned &= row_size == 0 || begin == first_begin;
182
4
            }
183
4
        }
184
185
        // Fast path: all visible rows already address the same physical nested positions.
186
2
        if (offsets_aligned) {
187
0
            return {};
188
0
        }
189
190
        // Slow path: build the shared offsets once, copy only selected row slices, and preserve
191
        // the outer row count by emitting repeated offsets for unselected rows.
192
2
        auto offsets_column = ColumnArray::ColumnOffsets::create();
193
2
        auto& normalized_offsets = offsets_column->get_data();
194
2
        normalized_offsets.reserve(batch_size);
195
2
        size_t normalized_offset = 0;
196
2
        const auto& first_offsets =
197
2
                assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0])
198
2
                        .get_offsets();
199
8
        for (size_t row = 0; row < batch_size; ++row) {
200
6
            if (is_selected(row)) {
201
4
                normalized_offset += first_offsets[row] - first_offsets[row - 1];
202
4
            }
203
6
            normalized_offsets.push_back(normalized_offset);
204
6
        }
205
2
        ColumnPtr shared_offsets = std::move(offsets_column);
206
207
2
        Columns normalized_columns(num_arguments);
208
6
        for (size_t i = 0; i < num_arguments; ++i) {
209
4
            const auto& array =
210
4
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
211
4
            auto nested = array.get_data().clone_empty();
212
16
            for (size_t row = 0; row < batch_size; ++row) {
213
12
                if (is_selected(row)) {
214
8
                    const auto& offsets = array.get_offsets();
215
8
                    const size_t begin = offsets[row - 1];
216
8
                    const size_t row_size = offsets[row] - begin;
217
8
                    nested->insert_range_from(array.get_data(), begin, row_size);
218
8
                }
219
12
            }
220
4
            ColumnPtr nested_column = std::move(nested);
221
4
            normalized_columns[i] = ColumnArray::create(nested_column, shared_offsets);
222
4
        }
223
2
        return normalized_columns;
224
2
    }
_ZNK5doris24AggregateFunctionForEach23normalize_batch_columnsIZNKS0_22add_batch_single_placeEmPcPPKNS_7IColumnERNS_5ArenaEEUlmE_EESt6vectorINS_3COWIS3_E13immutable_ptrIS3_EESaISE_EEmS6_OT_
Line
Count
Source
157
12
                                    IsSelected&& is_selected) const {
158
12
        bool offsets_aligned = true;
159
        // First pass: reject size mismatches before allocating, and detect whether the
160
        // original nested columns can use one shared element index.
161
24
        for (size_t row = 0; row < batch_size; ++row) {
162
12
            if (!is_selected(row)) {
163
0
                continue;
164
0
            }
165
166
12
            const auto& first_array =
167
12
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0]);
168
12
            const auto& first_offsets = first_array.get_offsets();
169
12
            const size_t first_begin = first_offsets[row - 1];
170
12
            const size_t row_size = first_offsets[row] - first_begin;
171
12
            for (size_t i = 1; i < num_arguments; ++i) {
172
0
                const auto& array =
173
0
                        assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
174
0
                const auto& offsets = array.get_offsets();
175
0
                const size_t begin = offsets[row - 1];
176
0
                if (offsets[row] - begin != row_size) {
177
0
                    throw Exception(ErrorCode::INTERNAL_ERROR,
178
0
                                    "Arrays passed to {} aggregate function have different sizes",
179
0
                                    get_name());
180
0
                }
181
0
                offsets_aligned &= row_size == 0 || begin == first_begin;
182
0
            }
183
12
        }
184
185
        // Fast path: all visible rows already address the same physical nested positions.
186
12
        if (offsets_aligned) {
187
12
            return {};
188
12
        }
189
190
        // Slow path: build the shared offsets once, copy only selected row slices, and preserve
191
        // the outer row count by emitting repeated offsets for unselected rows.
192
0
        auto offsets_column = ColumnArray::ColumnOffsets::create();
193
0
        auto& normalized_offsets = offsets_column->get_data();
194
0
        normalized_offsets.reserve(batch_size);
195
0
        size_t normalized_offset = 0;
196
0
        const auto& first_offsets =
197
0
                assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0])
198
0
                        .get_offsets();
199
0
        for (size_t row = 0; row < batch_size; ++row) {
200
0
            if (is_selected(row)) {
201
0
                normalized_offset += first_offsets[row] - first_offsets[row - 1];
202
0
            }
203
0
            normalized_offsets.push_back(normalized_offset);
204
0
        }
205
0
        ColumnPtr shared_offsets = std::move(offsets_column);
206
207
0
        Columns normalized_columns(num_arguments);
208
0
        for (size_t i = 0; i < num_arguments; ++i) {
209
0
            const auto& array =
210
0
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
211
0
            auto nested = array.get_data().clone_empty();
212
0
            for (size_t row = 0; row < batch_size; ++row) {
213
0
                if (is_selected(row)) {
214
0
                    const auto& offsets = array.get_offsets();
215
0
                    const size_t begin = offsets[row - 1];
216
0
                    const size_t row_size = offsets[row] - begin;
217
0
                    nested->insert_range_from(array.get_data(), begin, row_size);
218
0
                }
219
0
            }
220
0
            ColumnPtr nested_column = std::move(nested);
221
0
            normalized_columns[i] = ColumnArray::create(nested_column, shared_offsets);
222
0
        }
223
0
        return normalized_columns;
224
12
    }
aggregate_function_exception_test.cpp:_ZNK5doris24AggregateFunctionForEach23normalize_batch_columnsIZNS_79AggregateFunctionExceptionTest_ForEachNormalizesShiftedOffsetsOncePerBatch_Test8TestBodyEvE3$_0EESt6vectorINS_3COWINS_7IColumnEE13immutable_ptrIS6_EESaIS9_EEmPPKS6_OT_
Line
Count
Source
157
1
                                    IsSelected&& is_selected) const {
158
1
        bool offsets_aligned = true;
159
        // First pass: reject size mismatches before allocating, and detect whether the
160
        // original nested columns can use one shared element index.
161
4
        for (size_t row = 0; row < batch_size; ++row) {
162
3
            if (!is_selected(row)) {
163
1
                continue;
164
1
            }
165
166
2
            const auto& first_array =
167
2
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0]);
168
2
            const auto& first_offsets = first_array.get_offsets();
169
2
            const size_t first_begin = first_offsets[row - 1];
170
2
            const size_t row_size = first_offsets[row] - first_begin;
171
4
            for (size_t i = 1; i < num_arguments; ++i) {
172
2
                const auto& array =
173
2
                        assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
174
2
                const auto& offsets = array.get_offsets();
175
2
                const size_t begin = offsets[row - 1];
176
2
                if (offsets[row] - begin != row_size) {
177
0
                    throw Exception(ErrorCode::INTERNAL_ERROR,
178
0
                                    "Arrays passed to {} aggregate function have different sizes",
179
0
                                    get_name());
180
0
                }
181
2
                offsets_aligned &= row_size == 0 || begin == first_begin;
182
2
            }
183
2
        }
184
185
        // Fast path: all visible rows already address the same physical nested positions.
186
1
        if (offsets_aligned) {
187
0
            return {};
188
0
        }
189
190
        // Slow path: build the shared offsets once, copy only selected row slices, and preserve
191
        // the outer row count by emitting repeated offsets for unselected rows.
192
1
        auto offsets_column = ColumnArray::ColumnOffsets::create();
193
1
        auto& normalized_offsets = offsets_column->get_data();
194
1
        normalized_offsets.reserve(batch_size);
195
1
        size_t normalized_offset = 0;
196
1
        const auto& first_offsets =
197
1
                assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0])
198
1
                        .get_offsets();
199
4
        for (size_t row = 0; row < batch_size; ++row) {
200
3
            if (is_selected(row)) {
201
2
                normalized_offset += first_offsets[row] - first_offsets[row - 1];
202
2
            }
203
3
            normalized_offsets.push_back(normalized_offset);
204
3
        }
205
1
        ColumnPtr shared_offsets = std::move(offsets_column);
206
207
1
        Columns normalized_columns(num_arguments);
208
3
        for (size_t i = 0; i < num_arguments; ++i) {
209
2
            const auto& array =
210
2
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
211
2
            auto nested = array.get_data().clone_empty();
212
8
            for (size_t row = 0; row < batch_size; ++row) {
213
6
                if (is_selected(row)) {
214
4
                    const auto& offsets = array.get_offsets();
215
4
                    const size_t begin = offsets[row - 1];
216
4
                    const size_t row_size = offsets[row] - begin;
217
4
                    nested->insert_range_from(array.get_data(), begin, row_size);
218
4
                }
219
6
            }
220
2
            ColumnPtr nested_column = std::move(nested);
221
2
            normalized_columns[i] = ColumnArray::create(nested_column, shared_offsets);
222
2
        }
223
1
        return normalized_columns;
224
1
    }
_ZNK5doris24AggregateFunctionForEach23normalize_batch_columnsIRKZNKS_35AggregateFunctionNullVariadicInlineIS0_Lb1EE33streaming_agg_serialize_to_columnEPPKNS_7IColumnERNS_3COWIS4_E11mutable_ptrIS4_EEmRNS_5ArenaEEUlmE_EESt6vectorINS9_13immutable_ptrIS4_EESaISK_EEmS7_OT_
Line
Count
Source
157
2
                                    IsSelected&& is_selected) const {
158
2
        bool offsets_aligned = true;
159
        // First pass: reject size mismatches before allocating, and detect whether the
160
        // original nested columns can use one shared element index.
161
133
        for (size_t row = 0; row < batch_size; ++row) {
162
131
            if (!is_selected(row)) {
163
1
                continue;
164
1
            }
165
166
130
            const auto& first_array =
167
130
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0]);
168
130
            const auto& first_offsets = first_array.get_offsets();
169
130
            const size_t first_begin = first_offsets[row - 1];
170
130
            const size_t row_size = first_offsets[row] - first_begin;
171
260
            for (size_t i = 1; i < num_arguments; ++i) {
172
130
                const auto& array =
173
130
                        assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
174
130
                const auto& offsets = array.get_offsets();
175
130
                const size_t begin = offsets[row - 1];
176
130
                if (offsets[row] - begin != row_size) {
177
0
                    throw Exception(ErrorCode::INTERNAL_ERROR,
178
0
                                    "Arrays passed to {} aggregate function have different sizes",
179
0
                                    get_name());
180
0
                }
181
130
                offsets_aligned &= row_size == 0 || begin == first_begin;
182
130
            }
183
130
        }
184
185
        // Fast path: all visible rows already address the same physical nested positions.
186
2
        if (offsets_aligned) {
187
1
            return {};
188
1
        }
189
190
        // Slow path: build the shared offsets once, copy only selected row slices, and preserve
191
        // the outer row count by emitting repeated offsets for unselected rows.
192
1
        auto offsets_column = ColumnArray::ColumnOffsets::create();
193
1
        auto& normalized_offsets = offsets_column->get_data();
194
1
        normalized_offsets.reserve(batch_size);
195
1
        size_t normalized_offset = 0;
196
1
        const auto& first_offsets =
197
1
                assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0])
198
1
                        .get_offsets();
199
4
        for (size_t row = 0; row < batch_size; ++row) {
200
3
            if (is_selected(row)) {
201
2
                normalized_offset += first_offsets[row] - first_offsets[row - 1];
202
2
            }
203
3
            normalized_offsets.push_back(normalized_offset);
204
3
        }
205
1
        ColumnPtr shared_offsets = std::move(offsets_column);
206
207
1
        Columns normalized_columns(num_arguments);
208
3
        for (size_t i = 0; i < num_arguments; ++i) {
209
2
            const auto& array =
210
2
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
211
2
            auto nested = array.get_data().clone_empty();
212
8
            for (size_t row = 0; row < batch_size; ++row) {
213
6
                if (is_selected(row)) {
214
4
                    const auto& offsets = array.get_offsets();
215
4
                    const size_t begin = offsets[row - 1];
216
4
                    const size_t row_size = offsets[row] - begin;
217
4
                    nested->insert_range_from(array.get_data(), begin, row_size);
218
4
                }
219
6
            }
220
2
            ColumnPtr nested_column = std::move(nested);
221
2
            normalized_columns[i] = ColumnArray::create(nested_column, shared_offsets);
222
2
        }
223
1
        return normalized_columns;
224
2
    }
Unexecuted instantiation: _ZNK5doris24AggregateFunctionForEach23normalize_batch_columnsIRKZNKS_35AggregateFunctionNullVariadicInlineIS0_Lb0EE33streaming_agg_serialize_to_columnEPPKNS_7IColumnERNS_3COWIS4_E11mutable_ptrIS4_EEmRNS_5ArenaEEUlmE_EESt6vectorINS9_13immutable_ptrIS4_EESaISK_EEmS7_OT_
Unexecuted instantiation: _ZNK5doris24AggregateFunctionForEach23normalize_batch_columnsIRKZNKS_35AggregateFunctionNullVariadicInlineINS_26AggregateFunctionForEachV2ELb0EE33streaming_agg_serialize_to_columnEPPKNS_7IColumnERNS_3COWIS5_E11mutable_ptrIS5_EEmRNS_5ArenaEEUlmE_EESt6vectorINSA_13immutable_ptrIS5_EESaISL_EEmS8_OT_
Unexecuted instantiation: _ZNK5doris24AggregateFunctionForEach23normalize_batch_columnsIRKZNKS_35AggregateFunctionNullVariadicInlineINS_26AggregateFunctionForEachV2ELb1EE33streaming_agg_serialize_to_columnEPPKNS_7IColumnERNS_3COWIS5_E11mutable_ptrIS5_EEmRNS_5ArenaEEUlmE_EESt6vectorINSA_13immutable_ptrIS5_EESaISL_EEmS8_OT_
225
226
    template <typename IsSelected>
227
    PreparedBatchColumns prepare_batch_columns(size_t batch_size, const IColumn** columns,
228
14
                                               IsSelected&& is_selected) const {
229
14
        PreparedBatchColumns prepared;
230
14
        prepared.owned_columns =
231
14
                normalize_batch_columns(batch_size, columns, std::forward<IsSelected>(is_selected));
232
14
        prepared.nested_columns.resize(num_arguments);
233
30
        for (size_t i = 0; i < num_arguments; ++i) {
234
16
            const IColumn* column =
235
16
                    prepared.owned_columns.empty() ? columns[i] : prepared.owned_columns[i].get();
236
16
            const auto& array =
237
16
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*column);
238
16
            prepared.nested_columns[i] = &array.get_data();
239
16
            if (i == 0) {
240
14
                prepared.offsets = &array.get_offsets();
241
14
            }
242
16
        }
243
14
        return prepared;
244
14
    }
Unexecuted instantiation: _ZNK5doris24AggregateFunctionForEach21prepare_batch_columnsIZNKS0_9add_batchEmPPcmPPKNS_7IColumnERNS_5ArenaEbEUlmE_EENS0_20PreparedBatchColumnsEmS7_OT_
_ZNK5doris24AggregateFunctionForEach21prepare_batch_columnsIZNKS0_18add_batch_selectedEmPPcmPPKNS_7IColumnERNS_5ArenaEEUlmE_EENS0_20PreparedBatchColumnsEmS7_OT_
Line
Count
Source
228
2
                                               IsSelected&& is_selected) const {
229
2
        PreparedBatchColumns prepared;
230
2
        prepared.owned_columns =
231
2
                normalize_batch_columns(batch_size, columns, std::forward<IsSelected>(is_selected));
232
2
        prepared.nested_columns.resize(num_arguments);
233
6
        for (size_t i = 0; i < num_arguments; ++i) {
234
4
            const IColumn* column =
235
4
                    prepared.owned_columns.empty() ? columns[i] : prepared.owned_columns[i].get();
236
4
            const auto& array =
237
4
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*column);
238
4
            prepared.nested_columns[i] = &array.get_data();
239
4
            if (i == 0) {
240
2
                prepared.offsets = &array.get_offsets();
241
2
            }
242
4
        }
243
2
        return prepared;
244
2
    }
_ZNK5doris24AggregateFunctionForEach21prepare_batch_columnsIZNKS0_22add_batch_single_placeEmPcPPKNS_7IColumnERNS_5ArenaEEUlmE_EENS0_20PreparedBatchColumnsEmS6_OT_
Line
Count
Source
228
12
                                               IsSelected&& is_selected) const {
229
12
        PreparedBatchColumns prepared;
230
12
        prepared.owned_columns =
231
12
                normalize_batch_columns(batch_size, columns, std::forward<IsSelected>(is_selected));
232
12
        prepared.nested_columns.resize(num_arguments);
233
24
        for (size_t i = 0; i < num_arguments; ++i) {
234
12
            const IColumn* column =
235
12
                    prepared.owned_columns.empty() ? columns[i] : prepared.owned_columns[i].get();
236
12
            const auto& array =
237
12
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*column);
238
12
            prepared.nested_columns[i] = &array.get_data();
239
12
            if (i == 0) {
240
12
                prepared.offsets = &array.get_offsets();
241
12
            }
242
12
        }
243
12
        return prepared;
244
12
    }
245
246
public:
247
    constexpr static auto AGG_FOREACH_SUFFIX = "_foreach";
248
    AggregateFunctionForEach(AggregateFunctionPtr nested_function_, const DataTypes& arguments)
249
10
            : Base(arguments),
250
10
              nested_function {std::move(nested_function_)},
251
10
              nested_size_of_data(nested_function->size_of_data()),
252
10
              num_arguments(arguments.size()) {
253
10
        if (arguments.empty()) {
254
0
            throw Exception(ErrorCode::INTERNAL_ERROR,
255
0
                            "Aggregate function {} require at least one argument", get_name());
256
0
        }
257
10
    }
258
2
    void set_version(const int version_) override {
259
2
        Base::set_version(version_);
260
2
        nested_function->set_version(version_);
261
2
    }
262
263
2
    String get_name() const override { return nested_function->get_name() + AGG_FOREACH_SUFFIX; }
264
265
8
    DataTypePtr get_return_type() const override {
266
8
        return std::make_shared<DataTypeArray>(nested_function->get_return_type());
267
8
    }
268
269
    // Use foreach-specific streaming handling to align array inputs once while keeping aggregate
270
    // states row-local; prepare returns empty when the original columns are already aligned.
271
0
    bool requires_batch_add_for_streaming() const { return true; }
272
273
    template <typename IsSelected>
274
    Columns prepare_batch_columns_for_streaming(size_t batch_size, const IColumn** columns,
275
3
                                                IsSelected&& is_selected) const {
276
3
        return normalize_batch_columns(batch_size, columns, std::forward<IsSelected>(is_selected));
277
3
    }
aggregate_function_exception_test.cpp:_ZNK5doris24AggregateFunctionForEach35prepare_batch_columns_for_streamingIZNS_79AggregateFunctionExceptionTest_ForEachNormalizesShiftedOffsetsOncePerBatch_Test8TestBodyEvE3$_0EESt6vectorINS_3COWINS_7IColumnEE13immutable_ptrIS6_EESaIS9_EEmPPKS6_OT_
Line
Count
Source
275
1
                                                IsSelected&& is_selected) const {
276
1
        return normalize_batch_columns(batch_size, columns, std::forward<IsSelected>(is_selected));
277
1
    }
_ZNK5doris24AggregateFunctionForEach35prepare_batch_columns_for_streamingIRKZNKS_35AggregateFunctionNullVariadicInlineIS0_Lb1EE33streaming_agg_serialize_to_columnEPPKNS_7IColumnERNS_3COWIS4_E11mutable_ptrIS4_EEmRNS_5ArenaEEUlmE_EESt6vectorINS9_13immutable_ptrIS4_EESaISK_EEmS7_OT_
Line
Count
Source
275
2
                                                IsSelected&& is_selected) const {
276
2
        return normalize_batch_columns(batch_size, columns, std::forward<IsSelected>(is_selected));
277
2
    }
Unexecuted instantiation: _ZNK5doris24AggregateFunctionForEach35prepare_batch_columns_for_streamingIRKZNKS_35AggregateFunctionNullVariadicInlineIS0_Lb0EE33streaming_agg_serialize_to_columnEPPKNS_7IColumnERNS_3COWIS4_E11mutable_ptrIS4_EEmRNS_5ArenaEEUlmE_EESt6vectorINS9_13immutable_ptrIS4_EESaISK_EEmS7_OT_
Unexecuted instantiation: _ZNK5doris24AggregateFunctionForEach35prepare_batch_columns_for_streamingIRKZNKS_35AggregateFunctionNullVariadicInlineINS_26AggregateFunctionForEachV2ELb0EE33streaming_agg_serialize_to_columnEPPKNS_7IColumnERNS_3COWIS5_E11mutable_ptrIS5_EEmRNS_5ArenaEEUlmE_EESt6vectorINSA_13immutable_ptrIS5_EESaISL_EEmS8_OT_
Unexecuted instantiation: _ZNK5doris24AggregateFunctionForEach35prepare_batch_columns_for_streamingIRKZNKS_35AggregateFunctionNullVariadicInlineINS_26AggregateFunctionForEachV2ELb1EE33streaming_agg_serialize_to_columnEPPKNS_7IColumnERNS_3COWIS5_E11mutable_ptrIS5_EEmRNS_5ArenaEEUlmE_EESt6vectorINSA_13immutable_ptrIS5_EESaISL_EEmS8_OT_
278
279
172
    void destroy(AggregateDataPtr __restrict place) const noexcept override {
280
172
        AggregateFunctionForEachData& state = data(place);
281
282
172
        char* nested_state = state.array_of_aggregate_datas;
283
1.33k
        for (size_t i = 0; i < state.dynamic_array_size; ++i) {
284
1.16k
            nested_function->destroy(nested_state);
285
1.16k
            nested_state += nested_size_of_data;
286
1.16k
        }
287
172
    }
288
289
4
    bool is_trivial() const override {
290
4
        return std::is_trivial_v<Data> && nested_function->is_trivial();
291
4
    }
292
293
    void merge(AggregateDataPtr __restrict place, ConstAggregateDataPtr rhs,
294
9
               Arena& arena) const override {
295
9
        const AggregateFunctionForEachData& rhs_state = data(rhs);
296
9
        AggregateFunctionForEachData& state =
297
9
                ensure_aggregate_data(place, rhs_state.dynamic_array_size, arena);
298
299
9
        const char* rhs_nested_state = rhs_state.array_of_aggregate_datas;
300
9
        char* nested_state = state.array_of_aggregate_datas;
301
302
44
        for (size_t i = 0; i < state.dynamic_array_size && i < rhs_state.dynamic_array_size; ++i) {
303
35
            nested_function->merge(nested_state, rhs_nested_state, arena);
304
305
35
            rhs_nested_state += nested_size_of_data;
306
35
            nested_state += nested_size_of_data;
307
35
        }
308
9
    }
309
310
142
    void serialize(ConstAggregateDataPtr __restrict place, BufferWritable& buf) const override {
311
142
        const AggregateFunctionForEachData& state = data(place);
312
142
        buf.write_binary(state.dynamic_array_size);
313
142
        const char* nested_state = state.array_of_aggregate_datas;
314
1.21k
        for (size_t i = 0; i < state.dynamic_array_size; ++i) {
315
1.07k
            nested_function->serialize(nested_state, buf);
316
1.07k
            nested_state += nested_size_of_data;
317
1.07k
        }
318
142
    }
319
320
    void deserialize(AggregateDataPtr __restrict place, BufferReadable& buf,
321
11
                     Arena& arena) const override {
322
11
        AggregateFunctionForEachData& state = data(place);
323
324
11
        size_t new_size = 0;
325
11
        buf.read_binary(new_size);
326
327
11
        ensure_aggregate_data(place, new_size, arena);
328
329
11
        char* nested_state = state.array_of_aggregate_datas;
330
48
        for (size_t i = 0; i < new_size; ++i) {
331
37
            nested_function->deserialize(nested_state, buf, arena);
332
37
            nested_state += nested_size_of_data;
333
37
        }
334
11
    }
335
336
18
    void insert_result_into(ConstAggregateDataPtr __restrict place, IColumn& to) const override {
337
18
        const AggregateFunctionForEachData& state = data(place);
338
339
18
        auto& arr_to = assert_cast<ColumnArray&, TypeCheckOnRelease::DISABLE>(to);
340
18
        auto& offsets_to = arr_to.get_offsets();
341
18
        IColumn* elems_to = &arr_to.get_data();
342
18
        ColumnNullable* nullable_elems_to = nullptr;
343
18
        if (!nested_function->get_return_type()->is_nullable()) {
344
18
            nullable_elems_to = assert_cast<ColumnNullable*, TypeCheckOnRelease::DISABLE>(elems_to);
345
18
            elems_to = nullable_elems_to->get_nested_column_ptr().get();
346
18
        }
347
348
18
        char* nested_state = state.array_of_aggregate_datas;
349
75
        for (size_t i = 0; i < state.dynamic_array_size; ++i) {
350
57
            nested_function->insert_result_into(nested_state, *elems_to);
351
57
            if (nullable_elems_to != nullptr) {
352
57
                nullable_elems_to->get_null_map_data().push_back(0);
353
57
            }
354
57
            nested_state += nested_size_of_data;
355
57
        }
356
357
18
        offsets_to.push_back(offsets_to.back() + state.dynamic_array_size);
358
18
    }
359
360
12
    void check_result_column_type(const IColumn& to) const override {
361
12
        const auto* arr_to = check_and_get_column<ColumnArray>(to);
362
12
        if (UNLIKELY(arr_to == nullptr)) {
363
0
            throw doris::Exception(Status::InternalError(
364
0
                    "Aggregate function {} result type check failed: Column type {} is not "
365
0
                    "ColumnArray",
366
0
                    get_name(), to.get_name()));
367
0
        }
368
369
12
        const IColumn* elems_to = &arr_to->get_data();
370
12
        if (!nested_function->get_return_type()->is_nullable()) {
371
12
            const auto* nullable_elems_to = check_and_get_column<ColumnNullable>(*elems_to);
372
12
            if (UNLIKELY(nullable_elems_to == nullptr)) {
373
0
                throw doris::Exception(Status::InternalError(
374
0
                        "Aggregate function {} result type check failed: Array nested column "
375
0
                        "type {} is not ColumnNullable",
376
0
                        get_name(), elems_to->get_name()));
377
0
            }
378
12
            elems_to = &nullable_elems_to->get_nested_column();
379
12
        }
380
12
        nested_function->check_result_column_type(*elems_to);
381
12
    }
382
383
    void add(AggregateDataPtr __restrict place, const IColumn** columns, ssize_t row_num,
384
138
             Arena& arena) const override {
385
138
        std::vector<const IColumn*> nested(num_arguments);
386
387
407
        for (size_t i = 0; i < num_arguments; ++i) {
388
269
            nested[i] = &assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i])
389
269
                                 .get_data();
390
269
        }
391
392
138
        const auto& first_array_column =
393
138
                assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[0]);
394
138
        const auto& offsets = first_array_column.get_offsets();
395
396
138
        size_t begin = offsets[row_num - 1];
397
138
        size_t end = offsets[row_num];
398
138
        const size_t row_size = end - begin;
399
138
        bool offsets_aligned = true;
400
401
        /// Sanity check. NOTE We can implement specialization for a case with single argument, if the check will hurt performance.
402
269
        for (size_t i = 1; i < num_arguments; ++i) {
403
131
            const auto& ith_column =
404
131
                    assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
405
131
            const auto& ith_offsets = ith_column.get_offsets();
406
131
            const size_t ith_begin = ith_offsets[row_num - 1];
407
131
            const size_t ith_end = ith_offsets[row_num];
408
409
131
            if (ith_end - ith_begin != row_size) {
410
0
                throw Exception(ErrorCode::INTERNAL_ERROR,
411
0
                                "Arrays passed to {} aggregate function have different sizes",
412
0
                                get_name());
413
0
            }
414
131
            offsets_aligned &= ith_begin == begin;
415
131
        }
416
417
138
        std::vector<ColumnPtr> compacted_nested;
418
138
        if (!offsets_aligned && row_size != 0) {
419
            // The nested aggregate accepts one shared index, so align only mismatched row slices.
420
1
            compacted_nested.reserve(num_arguments);
421
3
            for (size_t i = 0; i < num_arguments; ++i) {
422
2
                const auto& ith_column =
423
2
                        assert_cast<const ColumnArray&, TypeCheckOnRelease::DISABLE>(*columns[i]);
424
2
                const size_t ith_begin = ith_column.get_offsets()[row_num - 1];
425
2
                compacted_nested.emplace_back(nested[i]->cut(ith_begin, row_size));
426
2
                nested[i] = compacted_nested.back().get();
427
2
            }
428
1
            begin = 0;
429
1
            end = row_size;
430
1
        }
431
432
138
        AggregateFunctionForEachData& state = ensure_aggregate_data(place, row_size, arena);
433
434
138
        char* nested_state = state.array_of_aggregate_datas;
435
1.18k
        for (size_t i = begin; i < end; ++i) {
436
1.04k
            nested_function->add(nested_state, nested.data(), i, arena);
437
1.04k
            nested_state += nested_size_of_data;
438
1.04k
        }
439
138
    }
440
441
    void add_batch(size_t batch_size, AggregateDataPtr* places, size_t place_offset,
442
0
                   const IColumn** columns, Arena& arena, bool /*agg_many*/) const override {
443
0
        auto prepared = prepare_batch_columns(batch_size, columns, [](size_t) { return true; });
444
0
        for (size_t row = 0; row < batch_size; ++row) {
445
0
            add_aligned(places[row] + place_offset, prepared.nested_columns.data(),
446
0
                        *prepared.offsets, row, arena);
447
0
        }
448
0
    }
449
450
    void add_batch_selected(size_t batch_size, AggregateDataPtr* places, size_t place_offset,
451
2
                            const IColumn** columns, Arena& arena) const override {
452
2
        auto prepared = prepare_batch_columns(batch_size, columns,
453
24
                                              [&](size_t row) { return places[row] != nullptr; });
454
8
        for (size_t row = 0; row < batch_size; ++row) {
455
6
            if (places[row] != nullptr) {
456
4
                add_aligned(places[row] + place_offset, prepared.nested_columns.data(),
457
4
                            *prepared.offsets, row, arena);
458
4
            }
459
6
        }
460
2
    }
461
462
    void add_batch_single_place(size_t batch_size, AggregateDataPtr place, const IColumn** columns,
463
12
                                Arena& arena) const override {
464
12
        auto prepared = prepare_batch_columns(batch_size, columns, [](size_t) { return true; });
465
24
        for (size_t row = 0; row < batch_size; ++row) {
466
12
            add_aligned(place, prepared.nested_columns.data(), *prepared.offsets, row, arena);
467
12
        }
468
12
    }
469
470
14
    void check_input_columns_type(const IColumn** columns) const override {
471
14
        std::vector<const IColumn*> nested(num_arguments);
472
28
        for (size_t i = 0; i < num_arguments; ++i) {
473
14
            const auto* array_column = check_and_get_column<ColumnArray>(*columns[i]);
474
14
            if (UNLIKELY(array_column == nullptr)) {
475
0
                throw doris::Exception(Status::InternalError(
476
0
                        "Aggregate function {} argument {} type check failed: Column type {} is "
477
0
                        "not ColumnArray",
478
0
                        get_name(), i, columns[i]->get_name()));
479
0
            }
480
14
            nested[i] = &array_column->get_data();
481
14
        }
482
14
        nested_function->check_input_columns_type(nested.data());
483
14
    }
484
};
485
} // namespace doris