Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion be/src/exprs/aggregate/aggregate_function_window_funnel.h
Original file line number Diff line number Diff line change
Expand Up @@ -193,7 +193,11 @@ struct WindowFunnelState {
if constexpr (T != TYPE_TIMESTAMP_NS) {
TimeInterval interval(SECOND, window, false);
end_timestamp = first_timestamp;
end_timestamp.template date_add_interval<SECOND>(interval);
if (!end_timestamp.template date_add_interval<SECOND>(interval)) {
throw Exception(ErrorCode::OUT_OF_BOUND,
"Operation window_funnel of {}, {} out of range",
first_timestamp.debug_string(), window);
}
}

matched_count++;
Expand Down
6 changes: 5 additions & 1 deletion be/src/exprs/aggregate/aggregate_function_window_funnel_v2.h
Original file line number Diff line number Diff line change
Expand Up @@ -261,7 +261,11 @@ struct WindowFunnelStateV2 {
} else {
DateValueType end_ts = _ts_from_int(base_ts);
TimeInterval interval(SECOND, window, false);
end_ts.template date_add_interval<SECOND>(interval);
if (!end_ts.template date_add_interval<SECOND>(interval)) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

建议改为比较实际微秒差

 可以复用已有的 datetime_diff_in_microseconds() (be/src/core/value/vdatetime_value.h:1274):

 const auto base = _ts_from_int(base_ts);
 const auto current = _ts_from_int(current_ts);
 return static_cast<__int128>(current.datetime_diff_in_microseconds(base)) <=
        static_cast<__int128>(window) * 1000000;

 这样不必构造可能超出日期上限的窗口终点,同时保留微秒精度。不能直接对 DATETIMEV2 的原始整数做减法,因为它是日期字段的位打包编码;也不宜用截断到整秒的差
 值,否则 1.000001 秒可能误入 1 秒窗口。

另外,V1 的日期加法路径 (be/src/exprs/aggregate/aggregate_function_window_funnel.h:194) 也忽略了同类返回值。现有 V2 单测包含纳秒上界用例,但未发现上述
DATETIMEV2 越界覆盖,建议补充该反例、恰好到窗口终点、超出 1 微秒及超大窗口用例。

throw Exception(ErrorCode::OUT_OF_BOUND,
"Operation window_funnel of {}, {} out of range",
end_ts.debug_string(), window);
}
return current_ts <= end_ts.to_date_int_val();
}
}
Expand Down
45 changes: 45 additions & 0 deletions be/test/exprs/aggregate/vec_window_funnel_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,51 @@ TEST_F(VWindowFunnelTest, testEmpty) {
agg_function->destroy(place2);
}

TEST_F(VWindowFunnelTest, testWindowOverflowThrows) {
AggregateFunctionSimpleFactory factory = AggregateFunctionSimpleFactory::instance();
DataTypes data_types = {std::make_shared<DataTypeInt64>(), std::make_shared<DataTypeString>(),
std::make_shared<DataTypeDateTimeV2>(),
std::make_shared<DataTypeUInt8>(), std::make_shared<DataTypeUInt8>()};
auto overflow_agg_function = factory.get("window_funnel", data_types, nullptr, false,
BeExecVersionManager::get_newest_version());
ASSERT_NE(overflow_agg_function, nullptr);

auto column_mode = ColumnString::create();
column_mode->insert(Field::create_field<TYPE_STRING>("default"));
column_mode->insert(Field::create_field<TYPE_STRING>("default"));

auto column_timestamp = ColumnDateTimeV2::create();
for (const auto& second : {58, 59}) {
VecDateTimeValue time_value;
time_value.unchecked_set_time(9999, 12, 31, 23, 59, second);
auto dtv2 = time_value.to_datetime_v2();
column_timestamp->insert_data((char*)&dtv2, 0);
}

auto column_window = ColumnInt64::create();
column_window->insert(Field::create_field<TYPE_BIGINT>(10));
column_window->insert(Field::create_field<TYPE_BIGINT>(10));
auto column_event1 = ColumnUInt8::create();
column_event1->insert(Field::create_field<TYPE_BOOLEAN>(1));
column_event1->insert(Field::create_field<TYPE_BOOLEAN>(0));
auto column_event2 = ColumnUInt8::create();
column_event2->insert(Field::create_field<TYPE_BOOLEAN>(0));
column_event2->insert(Field::create_field<TYPE_BOOLEAN>(1));

std::unique_ptr<char[]> memory(new char[overflow_agg_function->size_of_data()]);
AggregateDataPtr place = memory.get();
overflow_agg_function->create(place);
const IColumn* columns[] = {column_window.get(), column_mode.get(), column_timestamp.get(),
column_event1.get(), column_event2.get()};
for (int row = 0; row < 2; ++row) {
overflow_agg_function->add(place, columns, row, arena);
}

ColumnInt32 result;
EXPECT_THROW(overflow_agg_function->insert_result_into(place, result), Exception);
overflow_agg_function->destroy(place);
}

TEST_F(VWindowFunnelTest, testSerialize) {
const int NUM_CONDS = 4;
auto column_mode = ColumnString::create();
Expand Down
45 changes: 45 additions & 0 deletions be/test/exprs/aggregate/vec_window_funnel_v2_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,51 @@ TEST_F(VWindowFunnelV2Test, testEmpty) {
agg_function->destroy(place2);
}

TEST_F(VWindowFunnelV2Test, testWindowOverflowThrows) {
AggregateFunctionSimpleFactory factory = AggregateFunctionSimpleFactory::instance();
DataTypes data_types = {std::make_shared<DataTypeInt64>(), std::make_shared<DataTypeString>(),
std::make_shared<DataTypeDateTimeV2>(),
std::make_shared<DataTypeUInt8>(), std::make_shared<DataTypeUInt8>()};
auto overflow_agg_function = factory.get("window_funnel_v2", data_types, nullptr, false,
BeExecVersionManager::get_newest_version());
ASSERT_NE(overflow_agg_function, nullptr);

auto column_mode = ColumnString::create();
column_mode->insert(Field::create_field<TYPE_STRING>("default"));
column_mode->insert(Field::create_field<TYPE_STRING>("default"));

auto column_timestamp = ColumnDateTimeV2::create();
for (const auto& second : {58, 59}) {
VecDateTimeValue time_value;
time_value.unchecked_set_time(9999, 12, 31, 23, 59, second);
auto dtv2 = time_value.to_datetime_v2();
column_timestamp->insert_data((char*)&dtv2, 0);
}

auto column_window = ColumnInt64::create();
column_window->insert(Field::create_field<TYPE_BIGINT>(10));
column_window->insert(Field::create_field<TYPE_BIGINT>(10));
auto column_event1 = ColumnUInt8::create();
column_event1->insert(Field::create_field<TYPE_BOOLEAN>(1));
column_event1->insert(Field::create_field<TYPE_BOOLEAN>(0));
auto column_event2 = ColumnUInt8::create();
column_event2->insert(Field::create_field<TYPE_BOOLEAN>(0));
column_event2->insert(Field::create_field<TYPE_BOOLEAN>(1));

std::unique_ptr<char[]> memory(new char[overflow_agg_function->size_of_data()]);
AggregateDataPtr place = memory.get();
overflow_agg_function->create(place);
const IColumn* columns[] = {column_window.get(), column_mode.get(), column_timestamp.get(),
column_event1.get(), column_event2.get()};
for (int row = 0; row < 2; ++row) {
overflow_agg_function->add(place, columns, row, arena);
}

ColumnInt32 result;
EXPECT_THROW(overflow_agg_function->insert_result_into(place, result), Exception);
overflow_agg_function->destroy(place);
}

TEST_F(VWindowFunnelV2Test, testSerialize) {
const int NUM_CONDS = 4;
auto column_mode = ColumnString::create();
Expand Down
Loading