Coverage Report

Created: 2026-09-28 17:01

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exprs/aggregate/aggregate_function_window_funnel.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
18
// This file is copied from
19
// https://github.com/ClickHouse/ClickHouse/blob/master/AggregateFunctionWindowFunnel.h
20
// and modified by Doris
21
22
#pragma once
23
24
#include <gen_cpp/data.pb.h>
25
26
#include <algorithm>
27
#include <boost/iterator/iterator_facade.hpp>
28
#include <iterator>
29
#include <memory>
30
#include <type_traits>
31
#include <utility>
32
33
#include "common/cast_set.h"
34
#include "common/exception.h"
35
#include "common/status.h"
36
#include "core/assert_cast.h"
37
#include "core/binary_cast.hpp"
38
#include "core/column/column_string.h"
39
#include "core/column/column_vector.h"
40
#include "core/data_type/data_type_number.h"
41
#include "core/types.h"
42
#include "core/value/vdatetime_value.h"
43
#include "exec/sort/sort_block.h"
44
#include "exprs/aggregate/aggregate_function.h"
45
#include "util/simd/bits.h"
46
#include "util/var_int.h"
47
48
namespace doris {
49
class Arena;
50
class BufferReadable;
51
class BufferWritable;
52
class IColumn;
53
} // namespace doris
54
55
namespace doris {
56
57
enum class WindowFunnelMode : Int64 { INVALID, DEFAULT, DEDUPLICATION, FIXED, INCREASE };
58
59
1.67k
inline WindowFunnelMode string_to_window_funnel_mode(const String& string) {
60
1.67k
    if (string == "default") {
61
1.09k
        return WindowFunnelMode::DEFAULT;
62
1.09k
    } else if (string == "deduplication") {
63
72
        return WindowFunnelMode::DEDUPLICATION;
64
508
    } else if (string == "fixed") {
65
278
        return WindowFunnelMode::FIXED;
66
278
    } else if (string == "increase") {
67
36
        return WindowFunnelMode::INCREASE;
68
194
    } else {
69
194
        return WindowFunnelMode::INVALID;
70
194
    }
71
1.67k
}
72
73
template <PrimitiveType T>
74
struct DataValue {
75
    using TimestampEvent = std::vector<ColumnUInt8::Container>;
76
    using DateValueType = typename PrimitiveTypeTraits<T>::CppType;
77
    std::vector<DateValueType> dt;
78
    TimestampEvent event_columns_data;
79
    bool operator<(const DataValue& other) const { return dt < other.dt; }
80
264
    void clear() {
81
264
        dt.clear();
82
534
        for (auto& data : event_columns_data) {
83
534
            data.clear();
84
534
        }
85
264
    }
_ZN5doris9DataValueILNS_13PrimitiveTypeE43EE5clearEv
Line
Count
Source
80
2
    void clear() {
81
2
        dt.clear();
82
2
        for (auto& data : event_columns_data) {
83
2
            data.clear();
84
2
        }
85
2
    }
_ZN5doris9DataValueILNS_13PrimitiveTypeE26EE5clearEv
Line
Count
Source
80
262
    void clear() {
81
262
        dt.clear();
82
532
        for (auto& data : event_columns_data) {
83
532
            data.clear();
84
532
        }
85
262
    }
Unexecuted instantiation: _ZN5doris9DataValueILNS_13PrimitiveTypeE42EE5clearEv
86
1.55k
    auto size() const { return dt.size(); }
_ZNK5doris9DataValueILNS_13PrimitiveTypeE43EE4sizeEv
Line
Count
Source
86
6
    auto size() const { return dt.size(); }
_ZNK5doris9DataValueILNS_13PrimitiveTypeE26EE4sizeEv
Line
Count
Source
86
1.55k
    auto size() const { return dt.size(); }
Unexecuted instantiation: _ZNK5doris9DataValueILNS_13PrimitiveTypeE42EE4sizeEv
87
726
    bool empty() const { return dt.empty(); }
_ZNK5doris9DataValueILNS_13PrimitiveTypeE26EE5emptyEv
Line
Count
Source
87
726
    bool empty() const { return dt.empty(); }
Unexecuted instantiation: _ZNK5doris9DataValueILNS_13PrimitiveTypeE43EE5emptyEv
Unexecuted instantiation: _ZNK5doris9DataValueILNS_13PrimitiveTypeE42EE5emptyEv
88
    std::string debug_string() const {
89
        std::string result = "\n" + std::to_string(dt.size()) + " " +
90
                             std::to_string(event_columns_data[0].size()) + "\n";
91
        for (size_t i = 0; i < dt.size(); ++i) {
92
            result += dt[i].debug_string() + " ,";
93
            for (const auto& event : event_columns_data) {
94
                result += std::to_string(event[i]) + ",";
95
            }
96
            result += "\n";
97
        }
98
        return result;
99
    }
100
};
101
102
template <PrimitiveType T>
103
struct WindowFunnelState {
104
    static constexpr PrimitiveType PType = T;
105
    using NativeType = typename PrimitiveTypeTraits<T>::StorageFieldType;
106
    using DateValueType = typename PrimitiveTypeTraits<T>::CppType;
107
    int event_count = 0;
108
    int64_t window;
109
    bool enable_mode;
110
    WindowFunnelMode window_funnel_mode;
111
    DataValue<T> events_list;
112
113
1.06k
    WindowFunnelState() {
114
1.06k
        event_count = 0;
115
1.06k
        window = 0;
116
1.06k
        window_funnel_mode = WindowFunnelMode::INVALID;
117
1.06k
    }
_ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EEC2Ev
Line
Count
Source
113
6
    WindowFunnelState() {
114
6
        event_count = 0;
115
6
        window = 0;
116
6
        window_funnel_mode = WindowFunnelMode::INVALID;
117
6
    }
_ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EEC2Ev
Line
Count
Source
113
1.06k
    WindowFunnelState() {
114
1.06k
        event_count = 0;
115
1.06k
        window = 0;
116
1.06k
        window_funnel_mode = WindowFunnelMode::INVALID;
117
1.06k
    }
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EEC2Ev
118
1.06k
    WindowFunnelState(int arg_event_count) : WindowFunnelState() {
119
1.06k
        event_count = arg_event_count;
120
1.06k
        events_list.event_columns_data.resize(event_count);
121
1.06k
    }
_ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EEC2Ei
Line
Count
Source
118
6
    WindowFunnelState(int arg_event_count) : WindowFunnelState() {
119
6
        event_count = arg_event_count;
120
6
        events_list.event_columns_data.resize(event_count);
121
6
    }
_ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EEC2Ei
Line
Count
Source
118
1.06k
    WindowFunnelState(int arg_event_count) : WindowFunnelState() {
119
1.06k
        event_count = arg_event_count;
120
1.06k
        events_list.event_columns_data.resize(event_count);
121
1.06k
    }
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EEC2Ei
122
123
24
    void reset() {
124
24
        events_list.clear();
125
24
        window = 0;
126
24
        window_funnel_mode = WindowFunnelMode::INVALID;
127
24
    }
_ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE5resetEv
Line
Count
Source
123
24
    void reset() {
124
24
        events_list.clear();
125
24
        window = 0;
126
24
        window_funnel_mode = WindowFunnelMode::INVALID;
127
24
    }
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE5resetEv
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE5resetEv
128
129
748
    void add(const IColumn** arg_columns, ssize_t row_num, int64_t win, WindowFunnelMode mode) {
130
748
        window = win;
131
748
        window_funnel_mode = enable_mode ? mode : WindowFunnelMode::DEFAULT;
132
748
        events_list.dt.emplace_back(
133
748
                assert_cast<const typename PrimitiveTypeTraits<PType>::ColumnType&,
134
748
                            TypeCheckOnRelease::DISABLE>(*arg_columns[2])
135
748
                        .get_data()[row_num]);
136
2.58k
        for (int i = 0; i < event_count; i++) {
137
1.83k
            events_list.event_columns_data[i].emplace_back(
138
1.83k
                    assert_cast<const ColumnUInt8&, TypeCheckOnRelease::DISABLE>(
139
1.83k
                            *arg_columns[3 + i])
140
1.83k
                            .get_data()[row_num]);
141
1.83k
        }
142
748
    }
_ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE3addEPPKNS_7IColumnEllNS_16WindowFunnelModeE
Line
Count
Source
129
748
    void add(const IColumn** arg_columns, ssize_t row_num, int64_t win, WindowFunnelMode mode) {
130
748
        window = win;
131
748
        window_funnel_mode = enable_mode ? mode : WindowFunnelMode::DEFAULT;
132
748
        events_list.dt.emplace_back(
133
748
                assert_cast<const typename PrimitiveTypeTraits<PType>::ColumnType&,
134
748
                            TypeCheckOnRelease::DISABLE>(*arg_columns[2])
135
748
                        .get_data()[row_num]);
136
2.58k
        for (int i = 0; i < event_count; i++) {
137
1.83k
            events_list.event_columns_data[i].emplace_back(
138
1.83k
                    assert_cast<const ColumnUInt8&, TypeCheckOnRelease::DISABLE>(
139
1.83k
                            *arg_columns[3 + i])
140
1.83k
                            .get_data()[row_num]);
141
1.83k
        }
142
748
    }
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE3addEPPKNS_7IColumnEllNS_16WindowFunnelModeE
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE3addEPPKNS_7IColumnEllNS_16WindowFunnelModeE
143
144
    // todo: rethink thid sort method.
145
458
    void sort() {
146
458
        auto num = events_list.size();
147
458
        std::vector<size_t> indices(num);
148
458
        std::iota(indices.begin(), indices.end(), 0);
149
458
        std::sort(indices.begin(), indices.end(),
150
458
                  [this](size_t i1, size_t i2) { return events_list.dt[i1] < events_list.dt[i2]; });
_ZZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE4sortEvENKUlmmE_clEmm
Line
Count
Source
150
368
                  [this](size_t i1, size_t i2) { return events_list.dt[i1] < events_list.dt[i2]; });
Unexecuted instantiation: _ZZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE4sortEvENKUlmmE_clEmm
Unexecuted instantiation: _ZZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE4sortEvENKUlmmE_clEmm
151
152
1.47k
        auto reorder = [&indices, &num](auto& vec) {
153
1.47k
            std::decay_t<decltype(vec)> temp;
154
1.47k
            temp.resize(num);
155
3.80k
            for (auto i = 0; i < num; i++) {
156
2.33k
                temp[i] = vec[indices[i]];
157
2.33k
            }
158
1.47k
            std::swap(vec, temp);
159
1.47k
        };
_ZZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE4sortEvENKUlRT_E_clISt6vectorINS_11DateV2ValueINS_19DateTimeV2ValueTypeEEESaISA_EEEEDaS4_
Line
Count
Source
152
458
        auto reorder = [&indices, &num](auto& vec) {
153
458
            std::decay_t<decltype(vec)> temp;
154
458
            temp.resize(num);
155
1.11k
            for (auto i = 0; i < num; i++) {
156
660
                temp[i] = vec[indices[i]];
157
660
            }
158
458
            std::swap(vec, temp);
159
458
        };
_ZZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE4sortEvENKUlRT_E_clINS_8PODArrayIhLm4096ENS_9AllocatorILb0ELb0ELb0ENS_22DefaultMemoryAllocatorELb1EEELm16ELm15EEEEEDaS4_
Line
Count
Source
152
1.01k
        auto reorder = [&indices, &num](auto& vec) {
153
1.01k
            std::decay_t<decltype(vec)> temp;
154
1.01k
            temp.resize(num);
155
2.68k
            for (auto i = 0; i < num; i++) {
156
1.67k
                temp[i] = vec[indices[i]];
157
1.67k
            }
158
1.01k
            std::swap(vec, temp);
159
1.01k
        };
Unexecuted instantiation: _ZZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE4sortEvENKUlRT_E_clISt6vectorINS_16TimeStampNsValueESaIS8_EEEEDaS4_
Unexecuted instantiation: _ZZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE4sortEvENKUlRT_E_clINS_8PODArrayIhLm4096ENS_9AllocatorILb0ELb0ELb0ENS_22DefaultMemoryAllocatorELb1EEELm16ELm15EEEEEDaS4_
Unexecuted instantiation: _ZZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE4sortEvENKUlRT_E_clISt6vectorINS_16TimestampTzValueESaIS8_EEEEDaS4_
Unexecuted instantiation: _ZZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE4sortEvENKUlRT_E_clINS_8PODArrayIhLm4096ENS_9AllocatorILb0ELb0ELb0ENS_22DefaultMemoryAllocatorELb1EEELm16ELm15EEEEEDaS4_
160
161
458
        reorder(events_list.dt);
162
1.01k
        for (auto& inner_vec : events_list.event_columns_data) {
163
1.01k
            reorder(inner_vec);
164
1.01k
        }
165
458
    }
_ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE4sortEv
Line
Count
Source
145
458
    void sort() {
146
458
        auto num = events_list.size();
147
458
        std::vector<size_t> indices(num);
148
458
        std::iota(indices.begin(), indices.end(), 0);
149
458
        std::sort(indices.begin(), indices.end(),
150
458
                  [this](size_t i1, size_t i2) { return events_list.dt[i1] < events_list.dt[i2]; });
151
152
458
        auto reorder = [&indices, &num](auto& vec) {
153
458
            std::decay_t<decltype(vec)> temp;
154
458
            temp.resize(num);
155
458
            for (auto i = 0; i < num; i++) {
156
458
                temp[i] = vec[indices[i]];
157
458
            }
158
458
            std::swap(vec, temp);
159
458
        };
160
161
458
        reorder(events_list.dt);
162
1.01k
        for (auto& inner_vec : events_list.event_columns_data) {
163
1.01k
            reorder(inner_vec);
164
1.01k
        }
165
458
    }
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE4sortEv
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE4sortEv
166
167
    bool _within_window(const DateValueType& base_timestamp,
168
118
                        const DateValueType& current_timestamp) const {
169
118
        if constexpr (T == TYPE_TIMESTAMP_NS) {
170
2
            const auto elapsed_nanos = static_cast<__int128>(current_timestamp.epoch_nanos()) -
171
2
                                       base_timestamp.epoch_nanos();
172
2
            const auto window_nanos =
173
2
                    static_cast<__int128>(window) * TimeStampNsValue::NANOS_PER_SECOND;
174
2
            return elapsed_nanos <= window_nanos;
175
2
        }
176
0
        return static_cast<__int128>(current_timestamp.datetime_diff_in_microseconds(
177
118
                       base_timestamp)) <= static_cast<__int128>(window) * 1000000;
178
118
    }
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE14_within_windowERKNS_16TimeStampNsValueES5_
Line
Count
Source
168
2
                        const DateValueType& current_timestamp) const {
169
2
        if constexpr (T == TYPE_TIMESTAMP_NS) {
170
2
            const auto elapsed_nanos = static_cast<__int128>(current_timestamp.epoch_nanos()) -
171
2
                                       base_timestamp.epoch_nanos();
172
2
            const auto window_nanos =
173
2
                    static_cast<__int128>(window) * TimeStampNsValue::NANOS_PER_SECOND;
174
2
            return elapsed_nanos <= window_nanos;
175
2
        }
176
0
        return static_cast<__int128>(current_timestamp.datetime_diff_in_microseconds(
177
2
                       base_timestamp)) <= static_cast<__int128>(window) * 1000000;
178
2
    }
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE14_within_windowERKNS_11DateV2ValueINS_19DateTimeV2ValueTypeEEES7_
Line
Count
Source
168
116
                        const DateValueType& current_timestamp) const {
169
        if constexpr (T == TYPE_TIMESTAMP_NS) {
170
            const auto elapsed_nanos = static_cast<__int128>(current_timestamp.epoch_nanos()) -
171
                                       base_timestamp.epoch_nanos();
172
            const auto window_nanos =
173
                    static_cast<__int128>(window) * TimeStampNsValue::NANOS_PER_SECOND;
174
            return elapsed_nanos <= window_nanos;
175
        }
176
116
        return static_cast<__int128>(current_timestamp.datetime_diff_in_microseconds(
177
116
                       base_timestamp)) <= static_cast<__int128>(window) * 1000000;
178
116
    }
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE14_within_windowERKNS_16TimestampTzValueES5_
179
180
    template <WindowFunnelMode WINDOW_FUNNEL_MODE>
181
516
    int _match_event_list(size_t& start_row, size_t row_count) const {
182
516
        int matched_count = 0;
183
516
        if (window < 0) {
184
0
            throw Exception(ErrorCode::INVALID_ARGUMENT,
185
0
                            "the sliding time window must be a positive integer, but got: {}",
186
0
                            window);
187
0
        }
188
516
        int column_idx = 0;
189
424
        const auto& timestamp_data = events_list.dt;
190
424
        const auto& first_event_data = events_list.event_columns_data[column_idx].data();
191
424
        auto match_row = simd::find_one(first_event_data, start_row, row_count);
192
424
        start_row = match_row + 1;
193
516
        if (match_row < row_count) {
194
288
            auto prev_timestamp = timestamp_data[match_row];
195
288
            const auto first_timestamp = prev_timestamp;
196
288
            matched_count++;
197
288
            column_idx++;
198
288
            auto last_match_row = match_row;
199
288
            ++match_row;
200
372
            for (; column_idx < event_count && match_row < row_count; column_idx++, match_row++) {
201
152
                const auto& event_data = events_list.event_columns_data[column_idx];
202
152
                if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::FIXED) {
203
8
                    if (event_data[match_row] == 1) {
204
0
                        auto current_timestamp = timestamp_data[match_row];
205
0
                        if (_within_window(first_timestamp, current_timestamp)) {
206
0
                            matched_count++;
207
0
                            continue;
208
0
                        }
209
0
                    }
210
8
                    break;
211
8
                }
212
8
                match_row = simd::find_one(event_data.data(), match_row, row_count);
213
152
                if (match_row < row_count) {
214
112
                    auto current_timestamp = timestamp_data[match_row];
215
112
                    bool is_matched = _within_window(first_timestamp, current_timestamp);
216
112
                    if (is_matched) {
217
84
                        if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::INCREASE) {
218
0
                            is_matched = current_timestamp > prev_timestamp;
219
0
                        }
220
84
                    }
221
112
                    if (!is_matched) {
222
28
                        break;
223
28
                    }
224
84
                    if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::INCREASE) {
225
0
                        prev_timestamp = timestamp_data[match_row];
226
0
                    }
227
84
                    if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::DEDUPLICATION) {
228
0
                        bool is_dup = false;
229
0
                        if (match_row != last_match_row + 1) {
230
0
                            for (int tmp_column_idx = 0; tmp_column_idx < column_idx;
231
0
                                 tmp_column_idx++) {
232
0
                                const auto& tmp_event_data =
233
0
                                        events_list.event_columns_data[tmp_column_idx].data();
234
0
                                auto dup_match_row = simd::find_one(tmp_event_data,
235
0
                                                                    last_match_row + 1, match_row);
236
0
                                if (dup_match_row < match_row) {
237
0
                                    is_dup = true;
238
0
                                    break;
239
0
                                }
240
0
                            }
241
0
                        }
242
0
                        if (is_dup) {
243
0
                            break;
244
0
                        }
245
0
                        last_match_row = match_row;
246
0
                    }
247
0
                    matched_count++;
248
84
                } else {
249
40
                    break;
250
40
                }
251
152
            }
252
288
        }
253
100
        return matched_count;
254
92
    }
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE17_match_event_listILNS_16WindowFunnelModeE1EEEiRmm
Line
Count
Source
181
2
    int _match_event_list(size_t& start_row, size_t row_count) const {
182
2
        int matched_count = 0;
183
2
        if (window < 0) {
184
0
            throw Exception(ErrorCode::INVALID_ARGUMENT,
185
0
                            "the sliding time window must be a positive integer, but got: {}",
186
0
                            window);
187
0
        }
188
2
        int column_idx = 0;
189
2
        const auto& timestamp_data = events_list.dt;
190
2
        const auto& first_event_data = events_list.event_columns_data[column_idx].data();
191
2
        auto match_row = simd::find_one(first_event_data, start_row, row_count);
192
2
        start_row = match_row + 1;
193
2
        if (match_row < row_count) {
194
2
            auto prev_timestamp = timestamp_data[match_row];
195
2
            const auto first_timestamp = prev_timestamp;
196
2
            matched_count++;
197
2
            column_idx++;
198
2
            auto last_match_row = match_row;
199
2
            ++match_row;
200
4
            for (; column_idx < event_count && match_row < row_count; column_idx++, match_row++) {
201
2
                const auto& event_data = events_list.event_columns_data[column_idx];
202
                if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::FIXED) {
203
                    if (event_data[match_row] == 1) {
204
                        auto current_timestamp = timestamp_data[match_row];
205
                        if (_within_window(first_timestamp, current_timestamp)) {
206
                            matched_count++;
207
                            continue;
208
                        }
209
                    }
210
                    break;
211
                }
212
2
                match_row = simd::find_one(event_data.data(), match_row, row_count);
213
2
                if (match_row < row_count) {
214
2
                    auto current_timestamp = timestamp_data[match_row];
215
2
                    bool is_matched = _within_window(first_timestamp, current_timestamp);
216
2
                    if (is_matched) {
217
                        if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::INCREASE) {
218
                            is_matched = current_timestamp > prev_timestamp;
219
                        }
220
2
                    }
221
2
                    if (!is_matched) {
222
0
                        break;
223
0
                    }
224
                    if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::INCREASE) {
225
                        prev_timestamp = timestamp_data[match_row];
226
                    }
227
                    if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::DEDUPLICATION) {
228
                        bool is_dup = false;
229
                        if (match_row != last_match_row + 1) {
230
                            for (int tmp_column_idx = 0; tmp_column_idx < column_idx;
231
                                 tmp_column_idx++) {
232
                                const auto& tmp_event_data =
233
                                        events_list.event_columns_data[tmp_column_idx].data();
234
                                auto dup_match_row = simd::find_one(tmp_event_data,
235
                                                                    last_match_row + 1, match_row);
236
                                if (dup_match_row < match_row) {
237
                                    is_dup = true;
238
                                    break;
239
                                }
240
                            }
241
                        }
242
                        if (is_dup) {
243
                            break;
244
                        }
245
                        last_match_row = match_row;
246
                    }
247
2
                    matched_count++;
248
2
                } else {
249
0
                    break;
250
0
                }
251
2
            }
252
2
        }
253
2
        return matched_count;
254
2
    }
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE17_match_event_listILNS_16WindowFunnelModeE2EEEiRmm
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE17_match_event_listILNS_16WindowFunnelModeE3EEEiRmm
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE17_match_event_listILNS_16WindowFunnelModeE4EEEiRmm
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE17_match_event_listILNS_16WindowFunnelModeE1EEEiRmm
Line
Count
Source
181
422
    int _match_event_list(size_t& start_row, size_t row_count) const {
182
422
        int matched_count = 0;
183
422
        if (window < 0) {
184
0
            throw Exception(ErrorCode::INVALID_ARGUMENT,
185
0
                            "the sliding time window must be a positive integer, but got: {}",
186
0
                            window);
187
0
        }
188
422
        int column_idx = 0;
189
422
        const auto& timestamp_data = events_list.dt;
190
422
        const auto& first_event_data = events_list.event_columns_data[column_idx].data();
191
422
        auto match_row = simd::find_one(first_event_data, start_row, row_count);
192
422
        start_row = match_row + 1;
193
422
        if (match_row < row_count) {
194
222
            auto prev_timestamp = timestamp_data[match_row];
195
222
            const auto first_timestamp = prev_timestamp;
196
222
            matched_count++;
197
222
            column_idx++;
198
222
            auto last_match_row = match_row;
199
222
            ++match_row;
200
304
            for (; column_idx < event_count && match_row < row_count; column_idx++, match_row++) {
201
142
                const auto& event_data = events_list.event_columns_data[column_idx];
202
                if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::FIXED) {
203
                    if (event_data[match_row] == 1) {
204
                        auto current_timestamp = timestamp_data[match_row];
205
                        if (_within_window(first_timestamp, current_timestamp)) {
206
                            matched_count++;
207
                            continue;
208
                        }
209
                    }
210
                    break;
211
                }
212
142
                match_row = simd::find_one(event_data.data(), match_row, row_count);
213
142
                if (match_row < row_count) {
214
110
                    auto current_timestamp = timestamp_data[match_row];
215
110
                    bool is_matched = _within_window(first_timestamp, current_timestamp);
216
110
                    if (is_matched) {
217
                        if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::INCREASE) {
218
                            is_matched = current_timestamp > prev_timestamp;
219
                        }
220
82
                    }
221
110
                    if (!is_matched) {
222
28
                        break;
223
28
                    }
224
                    if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::INCREASE) {
225
                        prev_timestamp = timestamp_data[match_row];
226
                    }
227
                    if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::DEDUPLICATION) {
228
                        bool is_dup = false;
229
                        if (match_row != last_match_row + 1) {
230
                            for (int tmp_column_idx = 0; tmp_column_idx < column_idx;
231
                                 tmp_column_idx++) {
232
                                const auto& tmp_event_data =
233
                                        events_list.event_columns_data[tmp_column_idx].data();
234
                                auto dup_match_row = simd::find_one(tmp_event_data,
235
                                                                    last_match_row + 1, match_row);
236
                                if (dup_match_row < match_row) {
237
                                    is_dup = true;
238
                                    break;
239
                                }
240
                            }
241
                        }
242
                        if (is_dup) {
243
                            break;
244
                        }
245
                        last_match_row = match_row;
246
                    }
247
82
                    matched_count++;
248
82
                } else {
249
32
                    break;
250
32
                }
251
142
            }
252
222
        }
253
422
        return matched_count;
254
422
    }
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE17_match_event_listILNS_16WindowFunnelModeE2EEEiRmm
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE17_match_event_listILNS_16WindowFunnelModeE3EEEiRmm
Line
Count
Source
181
92
    int _match_event_list(size_t& start_row, size_t row_count) const {
182
92
        int matched_count = 0;
183
92
        if (window < 0) {
184
0
            throw Exception(ErrorCode::INVALID_ARGUMENT,
185
0
                            "the sliding time window must be a positive integer, but got: {}",
186
0
                            window);
187
0
        }
188
92
        int column_idx = 0;
189
92
        const auto& timestamp_data = events_list.dt;
190
92
        const auto& first_event_data = events_list.event_columns_data[column_idx].data();
191
92
        auto match_row = simd::find_one(first_event_data, start_row, row_count);
192
92
        start_row = match_row + 1;
193
92
        if (match_row < row_count) {
194
64
            auto prev_timestamp = timestamp_data[match_row];
195
64
            const auto first_timestamp = prev_timestamp;
196
64
            matched_count++;
197
64
            column_idx++;
198
64
            auto last_match_row = match_row;
199
64
            ++match_row;
200
64
            for (; column_idx < event_count && match_row < row_count; column_idx++, match_row++) {
201
8
                const auto& event_data = events_list.event_columns_data[column_idx];
202
8
                if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::FIXED) {
203
8
                    if (event_data[match_row] == 1) {
204
0
                        auto current_timestamp = timestamp_data[match_row];
205
0
                        if (_within_window(first_timestamp, current_timestamp)) {
206
0
                            matched_count++;
207
0
                            continue;
208
0
                        }
209
0
                    }
210
8
                    break;
211
8
                }
212
8
                match_row = simd::find_one(event_data.data(), match_row, row_count);
213
8
                if (match_row < row_count) {
214
0
                    auto current_timestamp = timestamp_data[match_row];
215
0
                    bool is_matched = _within_window(first_timestamp, current_timestamp);
216
0
                    if (is_matched) {
217
                        if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::INCREASE) {
218
                            is_matched = current_timestamp > prev_timestamp;
219
                        }
220
0
                    }
221
0
                    if (!is_matched) {
222
0
                        break;
223
0
                    }
224
                    if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::INCREASE) {
225
                        prev_timestamp = timestamp_data[match_row];
226
                    }
227
                    if constexpr (WINDOW_FUNNEL_MODE == WindowFunnelMode::DEDUPLICATION) {
228
                        bool is_dup = false;
229
                        if (match_row != last_match_row + 1) {
230
                            for (int tmp_column_idx = 0; tmp_column_idx < column_idx;
231
                                 tmp_column_idx++) {
232
                                const auto& tmp_event_data =
233
                                        events_list.event_columns_data[tmp_column_idx].data();
234
                                auto dup_match_row = simd::find_one(tmp_event_data,
235
                                                                    last_match_row + 1, match_row);
236
                                if (dup_match_row < match_row) {
237
                                    is_dup = true;
238
                                    break;
239
                                }
240
                            }
241
                        }
242
                        if (is_dup) {
243
                            break;
244
                        }
245
                        last_match_row = match_row;
246
                    }
247
0
                    matched_count++;
248
8
                } else {
249
8
                    break;
250
8
                }
251
8
            }
252
64
        }
253
100
        return matched_count;
254
92
    }
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE17_match_event_listILNS_16WindowFunnelModeE4EEEiRmm
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE17_match_event_listILNS_16WindowFunnelModeE1EEEiRmm
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE17_match_event_listILNS_16WindowFunnelModeE2EEEiRmm
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE17_match_event_listILNS_16WindowFunnelModeE3EEEiRmm
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE17_match_event_listILNS_16WindowFunnelModeE4EEEiRmm
255
256
    template <WindowFunnelMode WINDOW_FUNNEL_MODE>
257
448
    int _get_internal() const {
258
448
        size_t start_row = 0;
259
448
        int max_found_event_count = 0;
260
448
        auto row_count = events_list.size();
261
944
        while (start_row < row_count) {
262
516
            auto found_event_count = _match_event_list<WINDOW_FUNNEL_MODE>(start_row, row_count);
263
516
            if (found_event_count == event_count) {
264
20
                return found_event_count;
265
20
            }
266
496
            max_found_event_count = std::max(max_found_event_count, found_event_count);
267
496
        }
268
428
        return max_found_event_count;
269
448
    }
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE13_get_internalILNS_16WindowFunnelModeE1EEEiv
Line
Count
Source
257
2
    int _get_internal() const {
258
2
        size_t start_row = 0;
259
2
        int max_found_event_count = 0;
260
2
        auto row_count = events_list.size();
261
2
        while (start_row < row_count) {
262
2
            auto found_event_count = _match_event_list<WINDOW_FUNNEL_MODE>(start_row, row_count);
263
2
            if (found_event_count == event_count) {
264
2
                return found_event_count;
265
2
            }
266
0
            max_found_event_count = std::max(max_found_event_count, found_event_count);
267
0
        }
268
0
        return max_found_event_count;
269
2
    }
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE13_get_internalILNS_16WindowFunnelModeE2EEEiv
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE13_get_internalILNS_16WindowFunnelModeE3EEEiv
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE13_get_internalILNS_16WindowFunnelModeE4EEEiv
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE13_get_internalILNS_16WindowFunnelModeE1EEEiv
Line
Count
Source
257
362
    int _get_internal() const {
258
362
        size_t start_row = 0;
259
362
        int max_found_event_count = 0;
260
362
        auto row_count = events_list.size();
261
766
        while (start_row < row_count) {
262
422
            auto found_event_count = _match_event_list<WINDOW_FUNNEL_MODE>(start_row, row_count);
263
422
            if (found_event_count == event_count) {
264
18
                return found_event_count;
265
18
            }
266
404
            max_found_event_count = std::max(max_found_event_count, found_event_count);
267
404
        }
268
344
        return max_found_event_count;
269
362
    }
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE13_get_internalILNS_16WindowFunnelModeE2EEEiv
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE13_get_internalILNS_16WindowFunnelModeE3EEEiv
Line
Count
Source
257
84
    int _get_internal() const {
258
84
        size_t start_row = 0;
259
84
        int max_found_event_count = 0;
260
84
        auto row_count = events_list.size();
261
176
        while (start_row < row_count) {
262
92
            auto found_event_count = _match_event_list<WINDOW_FUNNEL_MODE>(start_row, row_count);
263
92
            if (found_event_count == event_count) {
264
0
                return found_event_count;
265
0
            }
266
92
            max_found_event_count = std::max(max_found_event_count, found_event_count);
267
92
        }
268
84
        return max_found_event_count;
269
84
    }
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE13_get_internalILNS_16WindowFunnelModeE4EEEiv
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE13_get_internalILNS_16WindowFunnelModeE1EEEiv
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE13_get_internalILNS_16WindowFunnelModeE2EEEiv
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE13_get_internalILNS_16WindowFunnelModeE3EEEiv
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE13_get_internalILNS_16WindowFunnelModeE4EEEiv
270
460
    int get() const {
271
460
        auto row_count = events_list.size();
272
460
        if (event_count == 0 || row_count == 0) {
273
12
            return 0;
274
12
        }
275
448
        switch (window_funnel_mode) {
276
364
        case WindowFunnelMode::DEFAULT:
277
364
            return _get_internal<WindowFunnelMode::DEFAULT>();
278
0
        case WindowFunnelMode::DEDUPLICATION:
279
0
            return _get_internal<WindowFunnelMode::DEDUPLICATION>();
280
84
        case WindowFunnelMode::FIXED:
281
84
            return _get_internal<WindowFunnelMode::FIXED>();
282
0
        case WindowFunnelMode::INCREASE:
283
0
            return _get_internal<WindowFunnelMode::INCREASE>();
284
0
        default:
285
0
            throw doris::Exception(ErrorCode::INTERNAL_ERROR, "Invalid window_funnel mode");
286
0
            return 0;
287
448
        }
288
448
    }
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE3getEv
Line
Count
Source
270
2
    int get() const {
271
2
        auto row_count = events_list.size();
272
2
        if (event_count == 0 || row_count == 0) {
273
0
            return 0;
274
0
        }
275
2
        switch (window_funnel_mode) {
276
2
        case WindowFunnelMode::DEFAULT:
277
2
            return _get_internal<WindowFunnelMode::DEFAULT>();
278
0
        case WindowFunnelMode::DEDUPLICATION:
279
0
            return _get_internal<WindowFunnelMode::DEDUPLICATION>();
280
0
        case WindowFunnelMode::FIXED:
281
0
            return _get_internal<WindowFunnelMode::FIXED>();
282
0
        case WindowFunnelMode::INCREASE:
283
0
            return _get_internal<WindowFunnelMode::INCREASE>();
284
0
        default:
285
0
            throw doris::Exception(ErrorCode::INTERNAL_ERROR, "Invalid window_funnel mode");
286
0
            return 0;
287
2
        }
288
2
    }
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE3getEv
Line
Count
Source
270
458
    int get() const {
271
458
        auto row_count = events_list.size();
272
458
        if (event_count == 0 || row_count == 0) {
273
12
            return 0;
274
12
        }
275
446
        switch (window_funnel_mode) {
276
362
        case WindowFunnelMode::DEFAULT:
277
362
            return _get_internal<WindowFunnelMode::DEFAULT>();
278
0
        case WindowFunnelMode::DEDUPLICATION:
279
0
            return _get_internal<WindowFunnelMode::DEDUPLICATION>();
280
84
        case WindowFunnelMode::FIXED:
281
84
            return _get_internal<WindowFunnelMode::FIXED>();
282
0
        case WindowFunnelMode::INCREASE:
283
0
            return _get_internal<WindowFunnelMode::INCREASE>();
284
0
        default:
285
0
            throw doris::Exception(ErrorCode::INTERNAL_ERROR, "Invalid window_funnel mode");
286
0
            return 0;
287
446
        }
288
446
    }
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE3getEv
289
290
418
    void merge(const WindowFunnelState<T>& other) {
291
418
        if (other.events_list.empty()) {
292
110
            return;
293
110
        }
294
295
308
        if (events_list.empty()) {
296
172
            window = other.window;
297
172
            window_funnel_mode = other.window_funnel_mode;
298
172
        } else if (UNLIKELY(window != other.window ||
299
136
                            window_funnel_mode != other.window_funnel_mode)) {
300
96
            throw Exception(ErrorCode::INVALID_ARGUMENT,
301
96
                            "window_funnel aggregate states have incompatible window or mode");
302
96
        }
303
212
        events_list.dt.insert(std::end(events_list.dt), std::begin(other.events_list.dt),
304
212
                              std::end(other.events_list.dt));
305
676
        for (size_t i = 0; i < event_count; i++) {
306
464
            events_list.event_columns_data[i].insert(
307
464
                    std::end(events_list.event_columns_data[i]),
308
464
                    std::begin(other.events_list.event_columns_data[i]),
309
464
                    std::end(other.events_list.event_columns_data[i]));
310
464
        }
311
212
        event_count = event_count > 0 ? event_count : other.event_count;
312
212
        window = window > 0 ? window : other.window;
313
212
        if (enable_mode) {
314
192
            window_funnel_mode = window_funnel_mode == WindowFunnelMode::INVALID
315
192
                                         ? other.window_funnel_mode
316
192
                                         : window_funnel_mode;
317
192
        } else {
318
20
            window_funnel_mode = WindowFunnelMode::DEFAULT;
319
20
        }
320
212
    }
_ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE5mergeERKS2_
Line
Count
Source
290
418
    void merge(const WindowFunnelState<T>& other) {
291
418
        if (other.events_list.empty()) {
292
110
            return;
293
110
        }
294
295
308
        if (events_list.empty()) {
296
172
            window = other.window;
297
172
            window_funnel_mode = other.window_funnel_mode;
298
172
        } else if (UNLIKELY(window != other.window ||
299
136
                            window_funnel_mode != other.window_funnel_mode)) {
300
96
            throw Exception(ErrorCode::INVALID_ARGUMENT,
301
96
                            "window_funnel aggregate states have incompatible window or mode");
302
96
        }
303
212
        events_list.dt.insert(std::end(events_list.dt), std::begin(other.events_list.dt),
304
212
                              std::end(other.events_list.dt));
305
676
        for (size_t i = 0; i < event_count; i++) {
306
464
            events_list.event_columns_data[i].insert(
307
464
                    std::end(events_list.event_columns_data[i]),
308
464
                    std::begin(other.events_list.event_columns_data[i]),
309
464
                    std::end(other.events_list.event_columns_data[i]));
310
464
        }
311
212
        event_count = event_count > 0 ? event_count : other.event_count;
312
212
        window = window > 0 ? window : other.window;
313
212
        if (enable_mode) {
314
192
            window_funnel_mode = window_funnel_mode == WindowFunnelMode::INVALID
315
192
                                         ? other.window_funnel_mode
316
192
                                         : window_funnel_mode;
317
192
        } else {
318
20
            window_funnel_mode = WindowFunnelMode::DEFAULT;
319
20
        }
320
212
    }
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE5mergeERKS2_
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE5mergeERKS2_
321
322
192
    void write(BufferWritable& out) const {
323
192
        write_var_int(event_count, out);
324
192
        write_var_int(window, out);
325
192
        if (enable_mode) {
326
188
            write_var_int(static_cast<std::underlying_type_t<WindowFunnelMode>>(window_funnel_mode),
327
188
                          out);
328
188
        }
329
192
        auto size = events_list.size();
330
192
        write_var_int(size, out);
331
192
        for (const auto& timestamp : events_list.dt) {
332
142
            write_var_int(timestamp.to_date_int_val(), out);
333
142
        }
334
582
        for (int64_t i = 0; i < event_count; i++) {
335
390
            const auto& event_columns_data = events_list.event_columns_data[i];
336
390
            for (auto event : event_columns_data) {
337
298
                write_var_int(event, out);
338
298
            }
339
390
        }
340
192
    }
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE5writeERNS_14BufferWritableE
Line
Count
Source
322
2
    void write(BufferWritable& out) const {
323
2
        write_var_int(event_count, out);
324
2
        write_var_int(window, out);
325
2
        if (enable_mode) {
326
2
            write_var_int(static_cast<std::underlying_type_t<WindowFunnelMode>>(window_funnel_mode),
327
2
                          out);
328
2
        }
329
2
        auto size = events_list.size();
330
2
        write_var_int(size, out);
331
2
        for (const auto& timestamp : events_list.dt) {
332
2
            write_var_int(timestamp.to_date_int_val(), out);
333
2
        }
334
4
        for (int64_t i = 0; i < event_count; i++) {
335
2
            const auto& event_columns_data = events_list.event_columns_data[i];
336
2
            for (auto event : event_columns_data) {
337
2
                write_var_int(event, out);
338
2
            }
339
2
        }
340
2
    }
_ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE5writeERNS_14BufferWritableE
Line
Count
Source
322
190
    void write(BufferWritable& out) const {
323
190
        write_var_int(event_count, out);
324
190
        write_var_int(window, out);
325
190
        if (enable_mode) {
326
186
            write_var_int(static_cast<std::underlying_type_t<WindowFunnelMode>>(window_funnel_mode),
327
186
                          out);
328
186
        }
329
190
        auto size = events_list.size();
330
190
        write_var_int(size, out);
331
190
        for (const auto& timestamp : events_list.dt) {
332
140
            write_var_int(timestamp.to_date_int_val(), out);
333
140
        }
334
578
        for (int64_t i = 0; i < event_count; i++) {
335
388
            const auto& event_columns_data = events_list.event_columns_data[i];
336
388
            for (auto event : event_columns_data) {
337
296
                write_var_int(event, out);
338
296
            }
339
388
        }
340
190
    }
Unexecuted instantiation: _ZNK5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE5writeERNS_14BufferWritableE
341
342
240
    void read(BufferReadable& in) {
343
240
        int64_t event_level;
344
240
        read_var_int(event_level, in);
345
240
        event_count = (int)event_level;
346
240
        read_var_int(window, in);
347
240
        window_funnel_mode = WindowFunnelMode::DEFAULT;
348
240
        if (enable_mode) {
349
236
            int64_t mode;
350
236
            read_var_int(mode, in);
351
236
            window_funnel_mode = static_cast<WindowFunnelMode>(mode);
352
236
        }
353
240
        int64_t size = 0;
354
240
        read_var_int(size, in);
355
240
        events_list.clear();
356
240
        events_list.dt.resize(size);
357
430
        for (auto i = 0; i < size; i++) {
358
190
            Int64 timestamp = 0;
359
190
            read_var_int(timestamp, in);
360
190
            if constexpr (T == TYPE_TIMESTAMP_NS) {
361
2
                events_list.dt[i] = DateValueType(timestamp);
362
188
            } else {
363
188
                events_list.dt[i] = DateValueType(static_cast<UInt64>(timestamp));
364
188
            }
365
190
        }
366
240
        events_list.event_columns_data.resize(event_count);
367
726
        for (int64_t i = 0; i < event_count; i++) {
368
486
            auto& event_columns_data = events_list.event_columns_data[i];
369
486
            event_columns_data.resize(size);
370
880
            for (auto j = 0; j < size; j++) {
371
394
                Int64 temp_value;
372
394
                read_var_int(temp_value, in);
373
394
                event_columns_data[j] = static_cast<UInt8>(temp_value);
374
394
            }
375
486
        }
376
240
    }
_ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE43EE4readERNS_14BufferReadableE
Line
Count
Source
342
2
    void read(BufferReadable& in) {
343
2
        int64_t event_level;
344
2
        read_var_int(event_level, in);
345
2
        event_count = (int)event_level;
346
2
        read_var_int(window, in);
347
2
        window_funnel_mode = WindowFunnelMode::DEFAULT;
348
2
        if (enable_mode) {
349
2
            int64_t mode;
350
2
            read_var_int(mode, in);
351
2
            window_funnel_mode = static_cast<WindowFunnelMode>(mode);
352
2
        }
353
2
        int64_t size = 0;
354
2
        read_var_int(size, in);
355
2
        events_list.clear();
356
2
        events_list.dt.resize(size);
357
4
        for (auto i = 0; i < size; i++) {
358
2
            Int64 timestamp = 0;
359
2
            read_var_int(timestamp, in);
360
2
            if constexpr (T == TYPE_TIMESTAMP_NS) {
361
2
                events_list.dt[i] = DateValueType(timestamp);
362
            } else {
363
                events_list.dt[i] = DateValueType(static_cast<UInt64>(timestamp));
364
            }
365
2
        }
366
2
        events_list.event_columns_data.resize(event_count);
367
4
        for (int64_t i = 0; i < event_count; i++) {
368
2
            auto& event_columns_data = events_list.event_columns_data[i];
369
2
            event_columns_data.resize(size);
370
4
            for (auto j = 0; j < size; j++) {
371
2
                Int64 temp_value;
372
2
                read_var_int(temp_value, in);
373
2
                event_columns_data[j] = static_cast<UInt8>(temp_value);
374
2
            }
375
2
        }
376
2
    }
_ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE26EE4readERNS_14BufferReadableE
Line
Count
Source
342
238
    void read(BufferReadable& in) {
343
238
        int64_t event_level;
344
238
        read_var_int(event_level, in);
345
238
        event_count = (int)event_level;
346
238
        read_var_int(window, in);
347
238
        window_funnel_mode = WindowFunnelMode::DEFAULT;
348
238
        if (enable_mode) {
349
234
            int64_t mode;
350
234
            read_var_int(mode, in);
351
234
            window_funnel_mode = static_cast<WindowFunnelMode>(mode);
352
234
        }
353
238
        int64_t size = 0;
354
238
        read_var_int(size, in);
355
238
        events_list.clear();
356
238
        events_list.dt.resize(size);
357
426
        for (auto i = 0; i < size; i++) {
358
188
            Int64 timestamp = 0;
359
188
            read_var_int(timestamp, in);
360
            if constexpr (T == TYPE_TIMESTAMP_NS) {
361
                events_list.dt[i] = DateValueType(timestamp);
362
188
            } else {
363
188
                events_list.dt[i] = DateValueType(static_cast<UInt64>(timestamp));
364
188
            }
365
188
        }
366
238
        events_list.event_columns_data.resize(event_count);
367
722
        for (int64_t i = 0; i < event_count; i++) {
368
484
            auto& event_columns_data = events_list.event_columns_data[i];
369
484
            event_columns_data.resize(size);
370
876
            for (auto j = 0; j < size; j++) {
371
392
                Int64 temp_value;
372
392
                read_var_int(temp_value, in);
373
392
                event_columns_data[j] = static_cast<UInt8>(temp_value);
374
392
            }
375
484
        }
376
238
    }
Unexecuted instantiation: _ZN5doris17WindowFunnelStateILNS_13PrimitiveTypeE42EE4readERNS_14BufferReadableE
377
};
378
379
template <PrimitiveType T>
380
class AggregateFunctionWindowFunnel final
381
        : public IAggregateFunctionDataHelper<WindowFunnelState<T>,
382
                                              AggregateFunctionWindowFunnel<T>>,
383
          MultiExpression,
384
          NullableAggregateFunction {
385
public:
386
    AggregateFunctionWindowFunnel(const DataTypes& argument_types_)
387
32
            : IAggregateFunctionDataHelper<WindowFunnelState<T>, AggregateFunctionWindowFunnel<T>>(
388
32
                      argument_types_) {}
_ZN5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EEC2ERKSt6vectorISt10shared_ptrIKNS_9IDataTypeEESaIS7_EE
Line
Count
Source
387
30
            : IAggregateFunctionDataHelper<WindowFunnelState<T>, AggregateFunctionWindowFunnel<T>>(
388
30
                      argument_types_) {}
_ZN5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EEC2ERKSt6vectorISt10shared_ptrIKNS_9IDataTypeEESaIS7_EE
Line
Count
Source
387
2
            : IAggregateFunctionDataHelper<WindowFunnelState<T>, AggregateFunctionWindowFunnel<T>>(
388
2
                      argument_types_) {}
Unexecuted instantiation: _ZN5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EEC2ERKSt6vectorISt10shared_ptrIKNS_9IDataTypeEESaIS7_EE
389
390
1.06k
    void create(AggregateDataPtr __restrict place) const override {
391
1.06k
        auto data = new (place) WindowFunnelState<T>(
392
1.06k
                cast_set<int>(IAggregateFunction::get_argument_types().size() - 3));
393
        /// support window funnel mode from 2.0. See `BeExecVersionManager::max_be_exec_version`
394
1.06k
        data->enable_mode = IAggregateFunction::version >= 3;
395
1.06k
    }
_ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EE6createEPc
Line
Count
Source
390
1.06k
    void create(AggregateDataPtr __restrict place) const override {
391
1.06k
        auto data = new (place) WindowFunnelState<T>(
392
1.06k
                cast_set<int>(IAggregateFunction::get_argument_types().size() - 3));
393
        /// support window funnel mode from 2.0. See `BeExecVersionManager::max_be_exec_version`
394
1.06k
        data->enable_mode = IAggregateFunction::version >= 3;
395
1.06k
    }
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EE6createEPc
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EE6createEPc
396
397
0
    String get_name() const override { return "window_funnel"; }
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EE8get_nameB5cxx11Ev
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EE8get_nameB5cxx11Ev
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EE8get_nameB5cxx11Ev
398
399
408
    DataTypePtr get_return_type() const override { return std::make_shared<DataTypeInt32>(); }
_ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EE15get_return_typeEv
Line
Count
Source
399
408
    DataTypePtr get_return_type() const override { return std::make_shared<DataTypeInt32>(); }
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EE15get_return_typeEv
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EE15get_return_typeEv
400
401
24
    void reset(AggregateDataPtr __restrict place) const override { this->data(place).reset(); }
_ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EE5resetEPc
Line
Count
Source
401
24
    void reset(AggregateDataPtr __restrict place) const override { this->data(place).reset(); }
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EE5resetEPc
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EE5resetEPc
402
403
    void add(AggregateDataPtr __restrict place, const IColumn** columns, ssize_t row_num,
404
748
             Arena&) const override {
405
748
        const auto& window =
406
748
                assert_cast<const ColumnInt64&, TypeCheckOnRelease::DISABLE>(*columns[0])
407
748
                        .get_data()[row_num];
408
748
        StringRef mode = columns[1]->get_data_at(row_num);
409
748
        this->data(place).add(columns, row_num, window,
410
748
                              string_to_window_funnel_mode(mode.to_string()));
411
748
    }
_ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EE3addEPcPPKNS_7IColumnElRNS_5ArenaE
Line
Count
Source
404
748
             Arena&) const override {
405
748
        const auto& window =
406
748
                assert_cast<const ColumnInt64&, TypeCheckOnRelease::DISABLE>(*columns[0])
407
748
                        .get_data()[row_num];
408
748
        StringRef mode = columns[1]->get_data_at(row_num);
409
748
        this->data(place).add(columns, row_num, window,
410
748
                              string_to_window_funnel_mode(mode.to_string()));
411
748
    }
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EE3addEPcPPKNS_7IColumnElRNS_5ArenaE
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EE3addEPcPPKNS_7IColumnElRNS_5ArenaE
412
413
0
    void check_input_columns_type(const IColumn** columns) const override {
414
0
        this->template check_argument_column_type<ColumnInt64>(columns[0]);
415
0
        this->template check_argument_column_type<ColumnString>(columns[1]);
416
0
        this->template check_argument_column_type<typename PrimitiveTypeTraits<T>::ColumnType>(
417
0
                columns[2]);
418
0
        for (size_t i = 3; i < this->argument_types.size(); ++i) {
419
0
            this->template check_argument_column_type<ColumnUInt8>(columns[i]);
420
0
        }
421
0
    }
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EE24check_input_columns_typeEPPKNS_7IColumnE
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EE24check_input_columns_typeEPPKNS_7IColumnE
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EE24check_input_columns_typeEPPKNS_7IColumnE
422
423
    void merge(AggregateDataPtr __restrict place, ConstAggregateDataPtr rhs,
424
418
               Arena&) const override {
425
418
        this->data(place).merge(this->data(rhs));
426
418
    }
_ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EE5mergeEPcPKcRNS_5ArenaE
Line
Count
Source
424
418
               Arena&) const override {
425
418
        this->data(place).merge(this->data(rhs));
426
418
    }
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EE5mergeEPcPKcRNS_5ArenaE
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EE5mergeEPcPKcRNS_5ArenaE
427
428
190
    void serialize(ConstAggregateDataPtr __restrict place, BufferWritable& buf) const override {
429
190
        this->data(place).write(buf);
430
190
    }
_ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EE9serializeEPKcRNS_14BufferWritableE
Line
Count
Source
428
190
    void serialize(ConstAggregateDataPtr __restrict place, BufferWritable& buf) const override {
429
190
        this->data(place).write(buf);
430
190
    }
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EE9serializeEPKcRNS_14BufferWritableE
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EE9serializeEPKcRNS_14BufferWritableE
431
432
    void deserialize(AggregateDataPtr __restrict place, BufferReadable& buf,
433
238
                     Arena&) const override {
434
238
        this->data(place).read(buf);
435
238
    }
_ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EE11deserializeEPcRNS_14BufferReadableERNS_5ArenaE
Line
Count
Source
433
238
                     Arena&) const override {
434
238
        this->data(place).read(buf);
435
238
    }
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EE11deserializeEPcRNS_14BufferReadableERNS_5ArenaE
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EE11deserializeEPcRNS_14BufferReadableERNS_5ArenaE
436
437
458
    void insert_result_into(ConstAggregateDataPtr __restrict place, IColumn& to) const override {
438
        // place is essentially an AggregateDataPtr, passed as a ConstAggregateDataPtr.
439
458
        this->data(const_cast<AggregateDataPtr>(place)).sort();
440
458
        assert_cast<ColumnInt32&, TypeCheckOnRelease::DISABLE>(to).get_data().push_back(
441
458
                IAggregateFunctionDataHelper<WindowFunnelState<T>,
442
458
                                             AggregateFunctionWindowFunnel<T>>::data(place)
443
458
                        .get());
444
458
    }
_ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE26EE18insert_result_intoEPKcRNS_7IColumnE
Line
Count
Source
437
458
    void insert_result_into(ConstAggregateDataPtr __restrict place, IColumn& to) const override {
438
        // place is essentially an AggregateDataPtr, passed as a ConstAggregateDataPtr.
439
458
        this->data(const_cast<AggregateDataPtr>(place)).sort();
440
458
        assert_cast<ColumnInt32&, TypeCheckOnRelease::DISABLE>(to).get_data().push_back(
441
458
                IAggregateFunctionDataHelper<WindowFunnelState<T>,
442
458
                                             AggregateFunctionWindowFunnel<T>>::data(place)
443
458
                        .get());
444
458
    }
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE43EE18insert_result_intoEPKcRNS_7IColumnE
Unexecuted instantiation: _ZNK5doris29AggregateFunctionWindowFunnelILNS_13PrimitiveTypeE42EE18insert_result_intoEPKcRNS_7IColumnE
445
446
protected:
447
    using IAggregateFunction::version;
448
};
449
} // namespace doris