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 |