asio多线程处理数据修复

This commit is contained in:
2026-06-29 10:01:33 +08:00
parent 8daf67aa5d
commit 8aac91872d
2 changed files with 10 additions and 4 deletions
+3 -4
View File
@@ -121,7 +121,6 @@ bool data_feed_debug = false;
asio::awaitable<void> Data_Feed::loop_coro() {
auto feed = this;
auto co = Coro::instance();
co_await asio::post(co->io, asio::use_awaitable);
auto& mode_acs = Global::instance()->mode_acs;
auto& cfg = mode_acs.data_feed_config;
auto& pool = cfg.pool_;
@@ -177,7 +176,7 @@ asio::awaitable<void> Data_Feed::loop_coro() {
std::cout << std::format("{} after send {}\n", key, thread_id_str());
}
}
co_await asio::post(co->io, asio::use_awaitable);
co_await resume_on(co->io.get_executor());
}
std::cout << std::format("{} feed loop exit {}\n", key, thread_id_str());
co_return;
@@ -186,7 +185,6 @@ bool data_source_debug = false;
asio::awaitable<void> Data_Source::loop_coro() {
auto source = this;
auto co = Coro::instance();
co_await asio::post(co->io, asio::use_awaitable);
while (source->running()) {
if (data_source_debug) {
std::cout << std::format("{} source loop begin {}\n", key, thread_id_str());
@@ -208,11 +206,12 @@ asio::awaitable<void> Data_Source::loop_coro() {
if (!enable || !source->running() || !source->registered()) {
break;
}
co_await asio::post(*co->process_data, asio::use_awaitable);
co_await resume_on(co->process_data->get_executor());
if (data_source_debug) {
std::cout << std::format("{} before process {}\n", key, thread_id_str());
}
auto num = source->process_mode_acs_data(mode_data);
co_await resume_on(co->io.get_executor());
if (data_source_debug) {
std::cout << std::format("{} after process {} {}\n", key, num, thread_id_str());
}
+7
View File
@@ -11,6 +11,7 @@
#include <iostream>
#include <memory>
#include <thread>
#include "asio/bind_executor.hpp"
#include "psc_global_include/Singleton.hpp"
class Coro : public Psc::Singleton<Coro> {
public:
@@ -32,3 +33,9 @@ private:
std::unique_ptr<std::future<void>> data_feed_thread_task;
std::atomic<bool> running = false;
};
inline asio::awaitable<void> resume_on(asio::any_io_executor ex) {
co_await asio::post(asio::bind_executor(ex, asio::deferred));
}
inline asio::awaitable<void> dispatch_on(asio::any_io_executor ex) {
co_await asio::dispatch(asio::bind_executor(ex, asio::deferred));
}