You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何为该回调场景实现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协程替代时遇到了两个核心问题:

  1. 如何向协程传入数据(比如Socket接收的数据);
  2. 处理数据时需要交还控制权给调用方,后续如何继续处理剩余数据?

问题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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 00:05:00