如何为该回调场景实现C++20协程替代方案?
用C++20协程替代回调方案的核心问题与解决思路
原有回调方案概述
我们有一套成熟的回调接口和数据处理类,伪代码如下:
回调接口定义
class ICallback { public: virtual int callbackA(T param1, U param2, V param3)=0; virtual int callbackB(int a, double b)=0; virtual int callbackC(T1 param1, T2 param2)=0; };
数据处理类
class Processor { public: explicit Processor(ICallback* callback) : m_callback(callback){} // 假设callback非空 int process(const uint8_t* data, int size) { int bytes_processed = 0; while(bytes_processed != size) { int bytes_consumed_in_this_iteration = doSomething(data + bytes_processed, size - bytes_processed); // doSomething可能通过m_callback调用任意回调方法 bytes_processed += bytes_consumed_in_this_iteration; } return size; // 返回值无关紧要 } private: int doSomething(const uint8_t* data, int size){ /* 实际处理逻辑 */ } ICallback* m_callback; };
这套代码功能正常,但尝试用C++20协程替代时遇到了两个核心问题:
- 如何向协程传入数据(比如Socket接收的数据);
- 处理数据时需要交还控制权给调用方,后续如何继续处理剩余数据?
问题1:向协程传入数据(以Socket为例)
核心思路是用线程安全的数据队列+可等待对象,让协程在没有数据时挂起,外部线程(如Socket接收线程)将数据推入队列后唤醒协程。
实现示例
// 封装待处理的数据块 struct DataChunk { const uint8_t* data; size_t size; }; // 协程等待数据的awaiter class DataAwaiter { public: explicit DataAwaiter(std::queue<DataChunk>& queue, std::mutex& mutex, std::condition_variable& cv) : m_queue(queue), m_mutex(mutex), m_cv(cv) {} // 检查是否已有数据,有则无需挂起 bool await_ready() const noexcept { std::lock_guard<std::mutex> lock(m_mutex); return !m_queue.empty(); } // 协程挂起时,等待数据并保存协程句柄 void await_suspend(std::coroutine_handle<> handle) noexcept { std::unique_lock<std::mutex> lock(m_mutex); m_waiting_handle = handle; // 等待数据或数据源关闭 m_cv.wait(lock, [this]() { return !m_queue.empty() || m_is_closed; }); } // 协程唤醒后,取出数据块返回 DataChunk await_resume() noexcept { std::lock_guard<std::mutex> lock(m_mutex); auto chunk = m_queue.front(); m_queue.pop(); return chunk; } // 外部调用:推入数据并唤醒协程 void push_data(const uint8_t* data, size_t size) { std::lock_guard<std::mutex> lock(m_mutex); m_queue.push({data, size}); m_cv.notify_one(); if (m_waiting_handle) { m_waiting_handle.resume(); m_waiting_handle = nullptr; } } // 关闭数据源(如Socket断开) void close() { std::lock_guard<std::mutex> lock(m_mutex); m_is_closed = true; m_cv.notify_one(); if (m_waiting_handle) { m_waiting_handle.resume(); m_waiting_handle = nullptr; } } private: std::queue<DataChunk>& m_queue; std::mutex& m_mutex; std::condition_variable& m_cv; std::coroutine_handle<> m_waiting_handle; bool m_is_closed = false; }; // 数据提供者:管理数据队列,给协程喂数据 class DataProvider { public: DataAwaiter wait_for_data() { return DataAwaiter(m_queue, m_mutex, m_cv); } void push(const uint8_t* data, size_t size) { std::lock_guard<std::mutex> lock(m_mutex); m_queue.push({data, size}); m_cv.notify_one(); } private: std::queue<DataChunk> m_queue; std::mutex m_mutex; std::condition_variable m_cv; };
协程使用示例
// 协程类型(已省略promise_type等样板代码) struct ProcessingTask { struct promise_type { ProcessingTask get_return_object() { return {}; } std::suspend_never initial_suspend() noexcept { return {}; } std::suspend_never final_suspend() noexcept { return {}; } void return_void() {} void unhandled_exception() { std::terminate(); } }; }; // 协程逻辑:持续等待并处理数据 ProcessingTask process_data(DataProvider& provider) { while (true) { auto chunk = co_await provider.wait_for_data(); if (chunk.size == 0) { // 数据源关闭,退出协程 co_return; } // 处理数据块(可替换为原有Processor的逻辑) process_chunk(chunk.data, chunk.size); } } // Socket接收线程逻辑:将接收到的数据推给协程 void socket_receive_thread(DataProvider& provider, int socket_fd) { uint8_t buffer[1024]; ssize_t bytes_read; while ((bytes_read = recv(socket_fd, buffer, sizeof(buffer), 0)) > 0) { provider.push(buffer, bytes_read); } // 通知协程数据源关闭 provider.push(nullptr, 0); }
问题2:处理数据时交还控制权,后续继续处理剩余数据
核心思路是把原有的批量处理拆成可中断的分步处理,每次处理一小段后主动挂起协程,调用方在合适时机恢复协程,继续处理剩余数据。
改造Processor为分步处理类
class StepwiseProcessor { public: StepwiseProcessor() = default; // 处理一部分数据,返回:已处理字节数 + 是否需要挂起(对应原回调触发时机) std::pair<size_t, bool> process_step(const uint8_t* data, size_t size) { size_t bytes_processed = 0; while (bytes_processed < size) { int consumed = doSomething(data + bytes_processed, size - bytes_processed); bytes_processed += consumed; // 模拟原回调触发时机:需要交还控制权 if (need_yield()) { // 保存当前剩余数据的状态(如果处理是有状态的,比如协议解析上下文) save_state(data + bytes_processed, size - bytes_processed); return {bytes_processed, true}; } } return {bytes_processed, false}; // 处理完整个数据块,无需挂起 } // 判断是否需要交还控制权(替换为实际业务逻辑) bool need_yield() { // 示例:每处理2字节就触发一次挂起 return m_processed_count % 2 == 0; } // 保存剩余数据状态 void save_state(const uint8_t* remaining_data, size_t remaining_size) { m_remaining_data = remaining_data; m_remaining_size = remaining_size; } // 获取剩余数据状态 std::pair<const uint8_t*, size_t> get_remaining_state() const { return {m_remaining_data, m_remaining_size}; } private: int doSomething(const uint8_t* data, int size) { m_processed_count++; return 1; // 示例:每次处理1字节 } bool m_need_yield = false; size_t m_processed_count = 0; const uint8_t* m_remaining_data = nullptr; size_t m_remaining_size = 0; };
协程中使用分步处理器
ProcessingTask run_processing(StepwiseProcessor& processor, DataProvider& provider) { while (true) { auto chunk = co_await provider.wait_for_data(); if (chunk.size == 0) break; const uint8_t* current_data = chunk.data; size_t current_size = chunk.size; while (current_size > 0) { auto [processed, need_yield] = processor.process_step(current_data, current_size); current_data += processed; current_size -= processed; if (need_yield) { // 主动挂起协程,交还控制权给调用方 co_await std::suspend_always{}; // 恢复后,取出剩余数据继续处理 auto [remaining_data, remaining_size] = processor.get_remaining_state(); current_data = remaining_data; current_size = remaining_size; } } } }
解释:当处理到需要触发原回调的逻辑时,协程通过co_await std::suspend_always{}主动挂起,控制权回到调用方。调用方处理完相关逻辑后,调用协程的resume()方法,协程会从挂起处继续处理剩余数据,代码全程是线性结构,避免了回调嵌套。
内容的提问来源于stack exchange,提问作者fyou
相关产品推荐
相关产品推荐

