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

如何实现主线程向长期运行子任务重复传值?Promise/Future遇阻

需求与问题

需要实现主线程向多个长期运行的子任务传递值,不允许通过终止重启子任务的方式传入新值。尝试用std::promise和std::future实现,但因为这两者不可重用,无法连续调用future.get(),导致代码无法正常运行。

以下是尝试的示例代码(存在问题):

#include <iostream>       // std::cout
#include <functional>     // std::ref
#include <thread>         // std::thread
#include <future>         // std::promise, std::future

void print_int(std::future<int>& input_future, std::promise<void>& done_promise)
{
    std::cout << "Print thread starts!" << std::endl;

    int x = input_future.get(); 
    while (x > 0)
    {
        std::cout << "Recived value: " << x;
        done_promise.set_value();
        std::cout << " done is set" << std::endl;

        int x = input_future.get();
    }
}


int main()
{
    std::cout << "Main thread starts!" << std::endl;

    std::promise<int> input_prom;                     
    std::promise<void> done_promise;
    std::future<void> done_future = done_promise.get_future();
    std::future<int> input_fut = input_prom.get_future();    
    std::thread thred(print_int, std::ref(input_fut), std::ref(done_promise)); 

    int counter = 10;

    while (counter > 1)
    {

        input_prom.set_value(counter);                         // fulfill promise

        //the print_int thread excecutes and prints

        done_future.get(); // wait until its printed
        counter--;
    }

    // (synchronizes with getting the future)
    thred.join();
    return 0;

}

问题分析

  1. std::promise只能被一次性设置值,调用set_value()后再次调用会直接抛出异常。
  2. std::future的get()方法同样只能调用一次,调用后future会进入无效状态,再次调用也会抛出异常。
  3. 示例中主线程循环调用input_prom.set_value(),以及子线程循环调用input_future.get(),都违反了这两个类的使用规则。

解决方案:条件变量+互斥锁+共享状态

要实现主线程持续向子线程传值且子线程长期运行,最常用的方案是条件变量+互斥锁搭配共享数据,或者使用线程安全的消息队列。下面是基于条件变量的实现:

实现代码

#include <iostream>
#include <thread>
#include <mutex>
#include <condition_variable>

// 共享数据结构,用于主线程和子线程间传递值
struct SharedData {
    std::mutex mtx;
    std::condition_variable cv;
    int value = 0;
    bool ready = false; // 标记是否有新值待处理
    bool stop = false;  // 标记是否终止子线程
};

void print_int(SharedData& data) {
    std::cout << "Print thread starts!" << std::endl;

    while (true) {
        std::unique_lock<std::mutex> lock(data.mtx);
        // 等待主线程发送新值或终止信号
        data.cv.wait(lock, [&data](){ return data.ready || data.stop; });

        if (data.stop) {
            break; // 收到终止信号,退出循环
        }

        // 处理收到的值
        int x = data.value;
        if (x <= 0) {
            break;
        }
        std::cout << "Received value: " << x << std::endl;

        // 标记值已处理完成
        data.ready = false;
        lock.unlock();
        data.cv.notify_one(); // 通知主线程已经处理完成
    }
    std::cout << "Print thread exits!" << std::endl;
}

int main() {
    std::cout << "Main thread starts!" << std::endl;

    SharedData data;
    std::thread thred(print_int, std::ref(data));

    int counter = 10;
    while (counter > 1) {
        std::unique_lock<std::mutex> lock(data.mtx);
        // 设置新值,并标记为待处理
        data.value = counter;
        data.ready = true;
        lock.unlock();
        data.cv.notify_one(); // 通知子线程有新值

        // 等待子线程处理完成
        std::unique_lock<std::mutex> wait_lock(data.mtx);
        data.cv.wait(wait_lock, [&data](){ return !data.ready; });

        counter--;
    }

    // 发送终止信号,让子线程退出
    std::unique_lock<std::mutex> lock(data.mtx);
    data.stop = true;
    lock.unlock();
    data.cv.notify_one();

    thred.join();
    std::cout << "Main thread exits!" << std::endl;
    return 0;
}

方案说明

  1. 共享状态管理:用SharedData结构体封装共享值、状态标记,配合互斥锁保证线程安全。
  2. 条件变量通知:主线程设置新值后通过cv.notify_one()唤醒子线程;子线程处理完成后同样通知主线程。
  3. 可重复使用:这种方案支持主线程持续传递新值,子线程无需重启,完全满足需求。
  4. 优雅终止:通过stop标记可以让子线程优雅退出,避免资源泄漏。

扩展:多子任务场景

如果需要给多个子线程传递值,只需要为每个子线程创建独立的SharedData实例,或者使用一个线程安全的消息队列(比如用std::queue配合互斥锁和条件变量),主线程向队列中放入值,各个子线程从队列中取值处理即可。

内容的提问来源于stack exchange,提问作者Suno

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 13:43:14