This is an automated email from the ASF dual-hosted git repository.
morningman pushed a commit to branch branch-4.0
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.0 by this push:
new 7910ec8cbec [fix](be) Avoid local runtime filter merge deadlock
(#65101)
7910ec8cbec is described below
commit 7910ec8cbec4afcbb0e97d67842c23afcb6440ce
Author: Pxl <[email protected]>
AuthorDate: Thu Jul 16 11:30:18 2026 +0800
[fix](be) Avoid local runtime filter merge deadlock (#65101)
### What problem does this PR solve?
Issue Number: None
Related PR: #64866
---
be/src/runtime/runtime_state.cpp | 2 +-
be/src/runtime_filter/runtime_filter_merger.h | 34 ++++--
be/src/runtime_filter/runtime_filter_mgr.cpp | 114 +++++++++++++--------
be/src/runtime_filter/runtime_filter_mgr.h | 14 +--
be/src/runtime_filter/runtime_filter_producer.cpp | 36 +++----
.../runtime_filter/runtime_filter_merger_test.cpp | 26 +++--
be/test/runtime_filter/runtime_filter_mgr_test.cpp | 34 +++---
7 files changed, 159 insertions(+), 101 deletions(-)
diff --git a/be/src/runtime/runtime_state.cpp b/be/src/runtime/runtime_state.cpp
index a99b9458972..dc1d19c0e42 100644
--- a/be/src/runtime/runtime_state.cpp
+++ b/be/src/runtime/runtime_state.cpp
@@ -494,7 +494,7 @@ Status RuntimeState::register_producer_runtime_filter(
// When RF is published, consumers in both global and local RF mgr will be
found.
RETURN_IF_ERROR(local_runtime_filter_mgr()->register_producer_filter(_query_ctx,
desc,
producer_filter));
-
RETURN_IF_ERROR(global_runtime_filter_mgr()->register_local_merger_producer_filter(
+
RETURN_IF_ERROR(global_runtime_filter_mgr()->register_local_merge_producer_filter(
_query_ctx, desc, *producer_filter));
return Status::OK();
}
diff --git a/be/src/runtime_filter/runtime_filter_merger.h
b/be/src/runtime_filter/runtime_filter_merger.h
index 00bac845473..166d31d50f5 100644
--- a/be/src/runtime_filter/runtime_filter_merger.h
+++ b/be/src/runtime_filter/runtime_filter_merger.h
@@ -17,6 +17,8 @@
#pragma once
+#include <algorithm>
+
#include "runtime_filter/runtime_filter.h"
#include "runtime_filter/runtime_filter_definitions.h"
#include "vec/exprs/vexpr.h"
@@ -47,6 +49,7 @@ public:
}
std::string debug_string() override {
+ std::unique_lock<std::recursive_mutex> l(_rmtx);
return fmt::format(
"Merger: ({}, expected_producer_num: {},
received_producer_num: {}, "
"received_rf_size_num: {}, received_sum_size: {})",
@@ -55,12 +58,15 @@ public:
}
// If input is a disabled predicate, the final result is a disabled
predicate.
- Status merge_from(const RuntimeFilter* other) {
+ // Returns true only for the call that makes the merger ready.
+ Status merge_from(const RuntimeFilter* other, bool* ready) {
+ std::unique_lock<std::recursive_mutex> l(_rmtx);
_received_producer_num++;
if (_expected_producer_num < _received_producer_num) {
return Status::InternalError(
"runtime filter merger input product more than expected,
{}", debug_string());
}
+ *ready = _received_producer_num == _expected_producer_num;
if (_received_producer_num == _expected_producer_num) {
_rf_state = State::READY;
}
@@ -72,13 +78,26 @@ public:
return st;
}
- void set_expected_producer_num(int num) {
- DCHECK_EQ(_received_producer_num, 0);
- DCHECK_EQ(_received_rf_size_num, 0);
- _expected_producer_num = num;
+ // Only raise the expected producer count. RuntimeFilterMgr may compute the
+ // count under its own lock and apply it after releasing that lock, so
+ // concurrent registrations can update the merger out of order.
+ void increase_expected_producer_num(int num) {
+ std::unique_lock<std::recursive_mutex> l(_rmtx);
+ if (_received_producer_num > 0 || _received_rf_size_num > 0) {
+ throw Exception(ErrorCode::INTERNAL_ERROR,
+ "runtime filter merger set expected producer after
receive data, {}",
+ debug_string());
+ }
+ _expected_producer_num = std::max(_expected_producer_num, num);
+ }
+
+ int get_expected_producer_num() {
+ std::unique_lock<std::recursive_mutex> l(_rmtx);
+ return _expected_producer_num;
}
bool add_rf_size(uint64_t size) {
+ std::unique_lock<std::recursive_mutex> l(_rmtx);
_received_rf_size_num++;
if (_expected_producer_num < _received_rf_size_num) {
throw Exception(ErrorCode::INTERNAL_ERROR,
@@ -89,7 +108,10 @@ public:
return (_received_rf_size_num == _expected_producer_num);
}
- uint64_t get_received_sum_size() const { return _received_sum_size; }
+ uint64_t get_received_sum_size() {
+ std::unique_lock<std::recursive_mutex> l(_rmtx);
+ return _received_sum_size;
+ }
bool ready() const { return _rf_state == State::READY; }
diff --git a/be/src/runtime_filter/runtime_filter_mgr.cpp
b/be/src/runtime_filter/runtime_filter_mgr.cpp
index b3c1cfa6c93..6debd1d0f66 100644
--- a/be/src/runtime_filter/runtime_filter_mgr.cpp
+++ b/be/src/runtime_filter/runtime_filter_mgr.cpp
@@ -67,13 +67,17 @@ Status RuntimeFilterMgr::register_consumer_filter(
SCOPED_CONSUME_MEM_TRACKER(_tracker.get());
int32_t key = desc.filter_id;
- std::lock_guard<std::mutex> l(_lock);
- RETURN_IF_ERROR(RuntimeFilterConsumer::create(query_ctx, &desc, node_id,
consumer));
- _consumer_map[key].push_back(*consumer);
+ std::shared_ptr<RuntimeFilterConsumer> new_consumer;
+ RETURN_IF_ERROR(RuntimeFilterConsumer::create(query_ctx, &desc, node_id,
&new_consumer));
+ {
+ std::lock_guard<std::mutex> l(_lock);
+ _consumer_map[key].push_back(new_consumer);
+ }
+ *consumer = new_consumer;
return Status::OK();
}
-Status RuntimeFilterMgr::register_local_merger_producer_filter(
+Status RuntimeFilterMgr::register_local_merge_producer_filter(
const QueryContext* query_ctx, const TRuntimeFilterDesc& desc,
std::shared_ptr<RuntimeFilterProducer> producer) {
if (!_is_global) [[unlikely]] {
@@ -89,44 +93,49 @@ Status
RuntimeFilterMgr::register_local_merger_producer_filter(
SCOPED_CONSUME_MEM_TRACKER(_tracker.get());
int32_t key = desc.filter_id;
- LocalMergeContext* context;
+ std::shared_ptr<LocalMergeContext> context;
+ std::shared_ptr<RuntimeFilterMerger> merger;
+ int expected_producer_num = 0;
{
std::lock_guard<std::mutex> l(_lock);
- context = &_local_merge_map[key]; // may inplace construct default
object
- }
-
- RETURN_IF_ERROR(context->register_producer(query_ctx, &desc, producer));
- return Status::OK();
-}
+ auto iter = _local_merge_map.find(key);
+ if (iter == _local_merge_map.end()) {
+ auto new_context = std::make_shared<LocalMergeContext>();
+ RETURN_IF_ERROR(RuntimeFilterMerger::create(query_ctx, &desc,
&new_context->merger));
+ _local_merge_map.emplace(key, new_context);
+ context = new_context;
+ } else {
+ DORIS_CHECK(iter->second);
+ context = iter->second;
+ }
-Status LocalMergeContext::register_producer(const QueryContext* query_ctx,
- const TRuntimeFilterDesc* desc,
-
std::shared_ptr<RuntimeFilterProducer> producer) {
- std::lock_guard<std::mutex> l(mtx);
- if (!merger) {
- RETURN_IF_ERROR(RuntimeFilterMerger::create(query_ctx, desc, &merger));
+ context->producers.emplace_back(producer);
+ merger = context->merger;
+ expected_producer_num = cast_set<int>(context->producers.size());
}
- producers.emplace_back(producer);
- merger->set_expected_producer_num(cast_set<int>(producers.size()));
+
+ merger->increase_expected_producer_num(expected_producer_num);
return Status::OK();
}
-Status RuntimeFilterMgr::get_local_merge_producer_filters(int filter_id,
- LocalMergeContext**
local_merge_filters) {
+Status RuntimeFilterMgr::get_local_merge_context(int filter_id,
+
std::shared_ptr<LocalMergeContext>* context) {
if (!_is_global) [[unlikely]] {
return Status::InternalError(
"A local merge filter can not be registered in Local
RuntimeFilterMgr");
}
+ context->reset();
std::lock_guard<std::mutex> l(_lock);
auto iter = _local_merge_map.find(filter_id);
if (iter == _local_merge_map.end()) {
return Status::InternalError(
- "get_local_merge_producer_filters meet unknown filter: {},
role: "
+ "get_local_merge_context meet unknown filter: {}, role: "
"LOCAL_MERGE_PRODUCER.",
filter_id);
}
- *local_merge_filters = &iter->second;
- DCHECK(iter->second.merger);
+ DORIS_CHECK(iter->second);
+ DORIS_CHECK(iter->second->merger);
+ *context = iter->second;
return Status::OK();
}
@@ -140,12 +149,22 @@ Status RuntimeFilterMgr::register_producer_filter(
SCOPED_CONSUME_MEM_TRACKER(_tracker.get());
int32_t key = desc.filter_id;
- std::lock_guard<std::mutex> l(_lock);
- if (_producer_id_set.contains(key)) {
- return Status::InvalidArgument("filter {} has been registered", key);
+ {
+ std::lock_guard<std::mutex> l(_lock);
+ if (_producer_id_set.contains(key)) {
+ return Status::InvalidArgument("filter {} has been registered",
key);
+ }
+ }
+ std::shared_ptr<RuntimeFilterProducer> new_producer;
+ RETURN_IF_ERROR(RuntimeFilterProducer::create(query_ctx, &desc,
&new_producer));
+ {
+ std::lock_guard<std::mutex> l(_lock);
+ if (_producer_id_set.contains(key)) {
+ return Status::InvalidArgument("filter {} has been registered",
key);
+ }
+ _producer_id_set.insert(key);
}
- RETURN_IF_ERROR(RuntimeFilterProducer::create(query_ctx, &desc, producer));
- _producer_id_set.insert(key);
+ *producer = new_producer;
return Status::OK();
}
@@ -184,7 +203,7 @@ Status RuntimeFilterMergeControllerEntity::_init_with_desc(
cnt_val->targetv2_info = targetv2_info;
RETURN_IF_ERROR(
RuntimeFilterMerger::create(query_ctx.get(), runtime_filter_desc,
&cnt_val->merger));
- cnt_val->merger->set_expected_producer_num(producer_size);
+ cnt_val->merger->increase_expected_producer_num(producer_size);
return Status::OK();
}
@@ -276,9 +295,9 @@ Status
RuntimeFilterMergeControllerEntity::send_filter_size(std::shared_ptr<Quer
}
Status RuntimeFilterMgr::sync_filter_size(const PSyncFilterSizeRequest*
request) {
- LocalMergeContext* local_merge_filters = nullptr;
- RETURN_IF_ERROR(get_local_merge_producer_filters(request->filter_id(),
&local_merge_filters));
- for (auto producer : local_merge_filters->producers) {
+ std::shared_ptr<LocalMergeContext> context;
+ RETURN_IF_ERROR(get_local_merge_context(request->filter_id(), &context));
+ for (const auto& producer : context->producers) {
producer->set_synced_size(request->filter_size());
}
return Status::OK();
@@ -286,18 +305,32 @@ Status RuntimeFilterMgr::sync_filter_size(const
PSyncFilterSizeRequest* request)
std::string RuntimeFilterMgr::debug_string() {
std::string result = "Local Merger Info:\n";
- std::lock_guard l(_lock);
- for (const auto& [filter_id, ctx] : _local_merge_map) {
+ struct LocalMergeContextSnapshot {
+ std::shared_ptr<RuntimeFilterMerger> merger;
+ std::vector<std::shared_ptr<RuntimeFilterProducer>> producers;
+ };
+ std::vector<LocalMergeContextSnapshot> local_merge_contexts;
+ std::vector<std::shared_ptr<RuntimeFilterConsumer>> consumers;
+ {
+ std::lock_guard l(_lock);
+ for (const auto& [filter_id, ctx] : _local_merge_map) {
+ DORIS_CHECK(ctx);
+ DORIS_CHECK(ctx->merger);
+ local_merge_contexts.push_back({ctx->merger, ctx->producers});
+ }
+ for (const auto& [filter_id, filter_consumers] : _consumer_map) {
+ consumers.insert(consumers.end(), filter_consumers.begin(),
filter_consumers.end());
+ }
+ }
+ for (const auto& ctx : local_merge_contexts) {
result += fmt::format("{}\n", ctx.merger->debug_string());
for (const auto& producer : ctx.producers) {
result += fmt::format("{}\n", producer->debug_string());
}
}
result += "Consumer Info:\n";
- for (const auto& [filter_id, consumers] : _consumer_map) {
- for (const auto& consumer : consumers) {
- result += fmt::format("{}\n", consumer->debug_string());
- }
+ for (const auto& consumer : consumers) {
+ result += fmt::format("{}\n", consumer->debug_string());
}
return result;
}
@@ -333,10 +366,9 @@ Status
RuntimeFilterMergeControllerEntity::merge(std::shared_ptr<QueryContext> q
RETURN_IF_ERROR(tmp_filter->assign(*request, attach_data));
- RETURN_IF_ERROR(cnt_val.merger->merge_from(tmp_filter.get()));
+ RETURN_IF_ERROR(cnt_val.merger->merge_from(tmp_filter.get(),
&is_ready));
cnt_val.arrive_id.insert(UniqueId(request->fragment_instance_id()));
- is_ready = cnt_val.merger->ready(); // update is_ready in locked scope
}
if (is_ready) {
diff --git a/be/src/runtime_filter/runtime_filter_mgr.h
b/be/src/runtime_filter/runtime_filter_mgr.h
index 2dcea3f7a4a..0680451ff21 100644
--- a/be/src/runtime_filter/runtime_filter_mgr.h
+++ b/be/src/runtime_filter/runtime_filter_mgr.h
@@ -57,12 +57,8 @@ class ExecEnv;
class RuntimeProfile;
struct LocalMergeContext {
- std::mutex mtx;
std::shared_ptr<RuntimeFilterMerger> merger;
std::vector<std::shared_ptr<RuntimeFilterProducer>> producers;
-
- Status register_producer(const QueryContext* query_ctx, const
TRuntimeFilterDesc* desc,
- std::shared_ptr<RuntimeFilterProducer> producer);
};
struct GlobalMergeContext {
@@ -87,11 +83,11 @@ public:
int node_id,
std::shared_ptr<RuntimeFilterConsumer>*
consumer_filter);
- Status register_local_merger_producer_filter(const QueryContext* query_ctx,
- const TRuntimeFilterDesc&
desc,
-
std::shared_ptr<RuntimeFilterProducer> producer);
+ Status register_local_merge_producer_filter(const QueryContext* query_ctx,
+ const TRuntimeFilterDesc& desc,
+
std::shared_ptr<RuntimeFilterProducer> producer);
- Status get_local_merge_producer_filters(int filter_id, LocalMergeContext**
local_merge_filters);
+ Status get_local_merge_context(int filter_id,
std::shared_ptr<LocalMergeContext>* context);
// Create local producer. This producer is hold by
RuntimeFilterProducerHelper.
Status register_producer_filter(const QueryContext* query_ctx, const
TRuntimeFilterDesc& desc,
@@ -122,7 +118,7 @@ private:
// key: "filter-id"
std::map<int32_t, std::vector<std::shared_ptr<RuntimeFilterConsumer>>>
_consumer_map;
std::set<int32_t> _producer_id_set;
- std::map<int32_t, LocalMergeContext> _local_merge_map;
+ std::map<int32_t, std::shared_ptr<LocalMergeContext>> _local_merge_map;
std::unique_ptr<MemTracker> _tracker;
diff --git a/be/src/runtime_filter/runtime_filter_producer.cpp
b/be/src/runtime_filter/runtime_filter_producer.cpp
index da5fbca15a9..fb14c2ade2e 100644
--- a/be/src/runtime_filter/runtime_filter_producer.cpp
+++ b/be/src/runtime_filter/runtime_filter_producer.cpp
@@ -54,12 +54,12 @@ Status RuntimeFilterProducer::publish(RuntimeState* state,
bool build_hash_table
// when global consumer not exist, send_to_local_targets will do
nothing, so merge rf is useless
return Status::OK();
}
- LocalMergeContext* context = nullptr;
-
RETURN_IF_ERROR(state->global_runtime_filter_mgr()->get_local_merge_producer_filters(
+ std::shared_ptr<LocalMergeContext> context;
+
RETURN_IF_ERROR(state->global_runtime_filter_mgr()->get_local_merge_context(
_wrapper->filter_id(), &context));
- std::lock_guard l(context->mtx);
- RETURN_IF_ERROR(context->merger->merge_from(this));
- if (context->merger->ready()) {
+ bool ready = false;
+ RETURN_IF_ERROR(context->merger->merge_from(this, &ready));
+ if (ready) {
if (_has_remote_target) {
RETURN_IF_ERROR(_send_to_remote_targets(state,
context->merger.get()));
} else {
@@ -157,22 +157,22 @@ Status RuntimeFilterProducer::send_size(RuntimeState*
state, uint64_t local_filt
set_state(State::WAITING_FOR_SYNCED_SIZE);
if (_need_do_merge(state)) {
- LocalMergeContext* merger_context = nullptr;
-
RETURN_IF_ERROR(state->global_runtime_filter_mgr()->get_local_merge_producer_filters(
- _wrapper->filter_id(), &merger_context));
- std::lock_guard merger_lock(merger_context->mtx);
- if (merger_context->merger->add_rf_size(local_filter_size)) {
- if (!_has_remote_target) {
- for (auto filter : merger_context->producers) {
-
filter->set_synced_size(merger_context->merger->get_received_sum_size());
- }
- return Status::OK();
- } else {
- local_filter_size =
merger_context->merger->get_received_sum_size();
+ std::shared_ptr<LocalMergeContext> context;
+
RETURN_IF_ERROR(state->global_runtime_filter_mgr()->get_local_merge_context(
+ _wrapper->filter_id(), &context));
+ uint64_t received_sum_size = 0;
+ bool ready_to_sync = context->merger->add_rf_size(local_filter_size);
+ if (!ready_to_sync) {
+ return Status::OK();
+ }
+ received_sum_size = context->merger->get_received_sum_size();
+ if (!_has_remote_target) {
+ for (const auto& filter : context->producers) {
+ filter->set_synced_size(received_sum_size);
}
- } else {
return Status::OK();
}
+ local_filter_size = received_sum_size;
} else if (!_has_remote_target) {
set_synced_size(local_filter_size);
diff --git a/be/test/runtime_filter/runtime_filter_merger_test.cpp
b/be/test/runtime_filter/runtime_filter_merger_test.cpp
index 2c62c0de8b0..0bc771fb1bf 100644
--- a/be/test/runtime_filter/runtime_filter_merger_test.cpp
+++ b/be/test/runtime_filter/runtime_filter_merger_test.cpp
@@ -35,15 +35,17 @@ public:
auto desc = TRuntimeFilterDescBuilder().build();
FAIL_IF_ERROR_OR_CATCH_EXCEPTION(
RuntimeFilterMerger::create(_query_ctx.get(), &desc, &merger));
- merger->set_expected_producer_num(2);
+ merger->increase_expected_producer_num(2);
ASSERT_FALSE(merger->ready());
ASSERT_EQ(merger->_wrapper->_state,
RuntimeFilterWrapper::State::UNINITED);
+ bool ready = false;
std::shared_ptr<RuntimeFilterProducer> producer;
FAIL_IF_ERROR_OR_CATCH_EXCEPTION(
_runtime_states[0]->register_producer_runtime_filter(desc,
&producer));
producer->set_wrapper_state_and_ready_to_publish(first_product_state);
- FAIL_IF_ERROR_OR_CATCH_EXCEPTION(merger->merge_from(producer.get()));
+ FAIL_IF_ERROR_OR_CATCH_EXCEPTION(merger->merge_from(producer.get(),
&ready));
+ ASSERT_FALSE(ready);
ASSERT_FALSE(merger->ready());
ASSERT_EQ(merger->_wrapper->_state, first_expected_state);
@@ -51,7 +53,8 @@ public:
FAIL_IF_ERROR_OR_CATCH_EXCEPTION(
_runtime_states[1]->register_producer_runtime_filter(desc,
&producer2));
producer2->set_wrapper_state_and_ready_to_publish(second_product_state);
- FAIL_IF_ERROR_OR_CATCH_EXCEPTION(merger->merge_from(producer2.get()));
+ FAIL_IF_ERROR_OR_CATCH_EXCEPTION(merger->merge_from(producer2.get(),
&ready));
+ ASSERT_TRUE(ready);
ASSERT_TRUE(merger->ready());
ASSERT_EQ(merger->_wrapper->_state, second_expected_state);
}
@@ -63,15 +66,17 @@ public:
std::shared_ptr<RuntimeFilterMerger> merger;
FAIL_IF_ERROR_OR_CATCH_EXCEPTION(
RuntimeFilterMerger::create(_query_ctx.get(), &desc, &merger));
- merger->set_expected_producer_num(1);
+ merger->increase_expected_producer_num(1);
ASSERT_FALSE(merger->ready());
+ bool ready = false;
std::shared_ptr<RuntimeFilterProducer> producer;
FAIL_IF_ERROR_OR_CATCH_EXCEPTION(
_runtime_states[0]->register_producer_runtime_filter(desc,
&producer));
FAIL_IF_ERROR_OR_CATCH_EXCEPTION(producer->init(123));
producer->set_wrapper_state_and_ready_to_publish(state);
- FAIL_IF_ERROR_OR_CATCH_EXCEPTION(merger->merge_from(producer.get()));
+ FAIL_IF_ERROR_OR_CATCH_EXCEPTION(merger->merge_from(producer.get(),
&ready));
+ ASSERT_TRUE(ready);
ASSERT_TRUE(merger->ready());
PMergeFilterRequest request;
@@ -99,7 +104,7 @@ TEST_F(RuntimeFilterMergerTest, add_rf_size) {
std::shared_ptr<RuntimeFilterMerger> merger;
auto desc = TRuntimeFilterDescBuilder().build();
FAIL_IF_ERROR_OR_CATCH_EXCEPTION(RuntimeFilterMerger::create(_query_ctx.get(),
&desc, &merger));
- merger->set_expected_producer_num(2);
+ merger->increase_expected_producer_num(2);
ASSERT_FALSE(merger->add_rf_size(123));
ASSERT_TRUE(merger->add_rf_size(1));
@@ -118,22 +123,25 @@ TEST_F(RuntimeFilterMergerTest, invalid_merge) {
std::shared_ptr<RuntimeFilterMerger> merger;
auto desc = TRuntimeFilterDescBuilder().build();
FAIL_IF_ERROR_OR_CATCH_EXCEPTION(RuntimeFilterMerger::create(_query_ctx.get(),
&desc, &merger));
- merger->set_expected_producer_num(1);
+ merger->increase_expected_producer_num(1);
ASSERT_FALSE(merger->ready());
ASSERT_EQ(merger->_wrapper->_state, RuntimeFilterWrapper::State::UNINITED);
+ bool ready = false;
std::shared_ptr<RuntimeFilterProducer> producer;
FAIL_IF_ERROR_OR_CATCH_EXCEPTION(
_runtime_states[0]->register_producer_runtime_filter(desc,
&producer));
producer->set_wrapper_state_and_ready_to_publish(RuntimeFilterWrapper::State::READY);
- FAIL_IF_ERROR_OR_CATCH_EXCEPTION(merger->merge_from(producer.get())); //
ready wrapper
+ FAIL_IF_ERROR_OR_CATCH_EXCEPTION(merger->merge_from(producer.get(),
&ready));
+ ASSERT_TRUE(ready);
+ ASSERT_TRUE(merger->ready());
ASSERT_EQ(merger->_wrapper->_state, RuntimeFilterWrapper::State::READY);
std::shared_ptr<RuntimeFilterProducer> producer2;
FAIL_IF_ERROR_OR_CATCH_EXCEPTION(
_runtime_states[1]->register_producer_runtime_filter(desc,
&producer2));
producer2->set_wrapper_state_and_ready_to_publish(RuntimeFilterWrapper::State::READY);
- auto st = merger->merge_from(producer2.get());
+ auto st = merger->merge_from(producer2.get(), &ready);
ASSERT_EQ(st.code(), ErrorCode::INTERNAL_ERROR);
}
diff --git a/be/test/runtime_filter/runtime_filter_mgr_test.cpp
b/be/test/runtime_filter/runtime_filter_mgr_test.cpp
index d8222e201d9..0f8a47bce75 100644
--- a/be/test/runtime_filter/runtime_filter_mgr_test.cpp
+++ b/be/test/runtime_filter/runtime_filter_mgr_test.cpp
@@ -77,12 +77,12 @@ TEST_F(RuntimeFilterMgrTest, TestRuntimeFilterMgr) {
// producer_filter should not be nullptr
EXPECT_FALSE(
global_runtime_filter_mgr
- ->register_local_merger_producer_filter(ctx.get(),
desc, producer_filter)
+ ->register_local_merge_producer_filter(ctx.get(),
desc, producer_filter)
.ok());
// local merge filter should not be registered in local mgr
EXPECT_FALSE(
local_runtime_filter_mgr
- ->register_local_merger_producer_filter(ctx.get(),
desc, producer_filter)
+ ->register_local_merge_producer_filter(ctx.get(),
desc, producer_filter)
.ok());
// producer should not registered in global mgr
EXPECT_FALSE(global_runtime_filter_mgr
@@ -103,24 +103,24 @@ TEST_F(RuntimeFilterMgrTest, TestRuntimeFilterMgr) {
.ok());
EXPECT_NE(producer_filter, nullptr);
- LocalMergeContext* local_merge_filters = nullptr;
- EXPECT_FALSE(global_runtime_filter_mgr
- ->get_local_merge_producer_filters(filter_id,
&local_merge_filters)
- .ok());
- EXPECT_FALSE(local_runtime_filter_mgr
- ->get_local_merge_producer_filters(filter_id,
&local_merge_filters)
- .ok());
- // Register local merge filter
- EXPECT_TRUE(
- global_runtime_filter_mgr
- ->register_local_merger_producer_filter(ctx.get(),
desc, producer_filter)
+ std::shared_ptr<LocalMergeContext> local_merge_context;
+ EXPECT_FALSE(
+ global_runtime_filter_mgr->get_local_merge_context(filter_id,
&local_merge_context)
+ .ok());
+ EXPECT_FALSE(
+ local_runtime_filter_mgr->get_local_merge_context(filter_id,
&local_merge_context)
.ok());
+ // Register local merge filter
EXPECT_TRUE(global_runtime_filter_mgr
- ->get_local_merge_producer_filters(filter_id,
&local_merge_filters)
+ ->register_local_merge_producer_filter(ctx.get(),
desc, producer_filter)
.ok());
- EXPECT_NE(local_merge_filters, nullptr);
- EXPECT_EQ(local_merge_filters->producers.size(), 1);
- local_merge_filters->producers.front()->_rf_state =
+ EXPECT_TRUE(
+ global_runtime_filter_mgr->get_local_merge_context(filter_id,
&local_merge_context)
+ .ok());
+ EXPECT_NE(local_merge_context, nullptr);
+ EXPECT_NE(local_merge_context->merger, nullptr);
+ EXPECT_EQ(local_merge_context->producers.size(), 1);
+ local_merge_context->producers.front()->_rf_state =
RuntimeFilterProducer::State ::WAITING_FOR_SYNCED_SIZE;
}
{
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]