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

C++11多线程回调定时器的线程安全实现及跨实例协调问题

C++11回调定时器实现与跨实例函数同步方案

问题描述

我需要实现一个回调定时器类,满足以下需求:

  • 构造时接收超时时间间隔和待执行函数两个参数
  • 每次超时后创建线程执行传入的函数,但该函数不保证线程安全
  • 若超时间隔短于函数执行时间,必须将回调请求排队,等前一次调用完成后再执行,禁止同一函数并发运行

另外有个疑问:如果主线程创建多个定时器实例,且这些实例传入同一个函数,能否跨实例协调对该函数的访问?

示例回调函数:

void myCallback() {
   // 执行业务操作
   std::cout << "回调函数执行中。\n";
   std::stringstream ss;

   int seconds = 5;
   for (int i = seconds; i > 0; i--)
   {
      ss << "线程ID: " << std::this_thread::get_id() << " 正在处理,剩余 " << i << " 秒\n";
      std::cout << ss.str();
      std::this_thread::sleep_for(std::chrono::seconds(1));

      ss.str(std::string());
   }
   ss << "线程ID: " << std::this_thread::get_id() << " 执行完成!\n";
   std::cout << ss.str();
}

一、单个定时器类的实现逻辑

核心思路是用任务队列缓存回调请求,配合互斥锁和条件变量实现串行执行,同时用后台线程处理定时调度。

关键设计点

  1. 线程安全的任务队列:用std::queue存储待执行的回调,通过std::mutex保护队列的读写操作
  2. 串行执行控制:用一个标记变量记录当前是否有回调在执行,确保同一时间只有一个线程运行目标函数
  3. 定时调度:后台线程循环等待超时,超时后将回调加入队列,同时监听队列状态处理任务
  4. 优雅退出:析构时标记停止状态,唤醒后台线程并等待其结束,避免资源泄漏

完整实现代码

#include <iostream>
#include <thread>
#include <chrono>
#include <mutex>
#include <condition_variable>
#include <queue>
#include <function>
#include <atomic>
#include <sstream>
#include <memory>

void myCallback() {
   // 执行业务操作
   std::cout << "回调函数执行中。\n";
   std::stringstream ss;

   int seconds = 5;
   for (int i = seconds; i > 0; i--)
   {
      ss << "线程ID: " << std::this_thread::get_id() << " 正在处理,剩余 " << i << " 秒\n";
      std::cout << ss.str();
      std::this_thread::sleep_for(std::chrono::seconds(1));

      ss.str(std::string());
   }
   ss << "线程ID: " << std::this_thread::get_id() << " 执行完成!\n";
   std::cout << ss.str();
}

class CallbackTimer {
public:
    CallbackTimer(std::chrono::milliseconds interval, std::function<void()> callback)
        : m_interval(interval), m_callback(std::move(callback)), m_running(true), m_is_executing(false) {
        m_thread = std::thread(&CallbackTimer::timerLoop, this);
    }

    ~CallbackTimer() {
        {
            std::lock_guard<std::mutex> lock(m_mutex);
            m_running = false;
        }
        m_cv.notify_all();
        if (m_thread.joinable()) {
            m_thread.join();
        }
    }

private:
    void timerLoop() {
        while (m_running) {
            // 等待超时或停止信号
            std::unique_lock<std::mutex> lock(m_mutex);
            bool stop_requested = m_cv.wait_for(lock, m_interval, [this]() { return !m_running; });
            
            if (stop_requested) {
                break;
            }

            // 将回调加入任务队列
            m_task_queue.push(m_callback);
            lock.unlock();
            m_cv.notify_all();

            // 尝试处理队列中的任务
            processTasks();
        }

        // 退出前处理剩余任务
        processTasks();
    }

    void processTasks() {
        std::function<void()> task;
        {
            std::lock_guard<std::mutex> lock(m_mutex);
            if (m_task_queue.empty() || m_is_executing) {
                return;
            }
            task = std::move(m_task_queue.front());
            m_task_queue.pop();
            m_is_executing = true;
        }

        if (task) {
            try {
                task();
            } catch (...) {
                std::cerr << "回调函数执行抛出异常\n";
            }
            std::lock_guard<std::mutex> lock(m_mutex);
            m_is_executing = false;
            m_cv.notify_all();
        }
    }

    std::chrono::milliseconds m_interval;
    std::function<void()> m_callback;
    std::queue<std::function<void()>> m_task_queue;
    std::mutex m_mutex;
    std::condition_variable m_cv;
    std::thread m_thread;
    std::atomic<bool> m_running;
    bool m_is_executing; // 受m_mutex保护
};

int main() {
    // 测试:2秒间隔的定时器
    CallbackTimer timer(std::chrono::seconds(2), myCallback);
    std::this_thread::sleep_for(std::chrono::seconds(15));
    return 0;
}

二、跨实例协调同一函数的访问

默认实现只能保证同一定时器实例内的回调串行执行,不同实例共享同一个函数时,需要额外的全局同步机制。

实现方案:共享互斥锁

让所有使用同一函数的定时器实例共享同一个互斥锁,在执行回调前加锁,确保全局串行。

修改后的定时器类

class CallbackTimer {
public:
    // 重载构造函数,支持传入共享互斥锁
    CallbackTimer(std::chrono::milliseconds interval, std::function<void()> callback, std::shared_ptr<std::mutex> shared_mutex = nullptr)
        : m_interval(interval), 
          m_callback(std::move(callback)), 
          m_shared_mutex(shared_mutex ? shared_mutex : std::make_shared<std::mutex>()),
          m_running(true), 
          m_is_executing(false) {
        m_thread = std::thread(&CallbackTimer::timerLoop, this);
    }

    ~CallbackTimer() {
        {
            std::lock_guard<std::mutex> lock(m_mutex);
            m_running = false;
        }
        m_cv.notify_all();
        if (m_thread.joinable()) {
            m_thread.join();
        }
    }

private:
    void timerLoop() {
        while (m_running) {
            std::unique_lock<std::mutex> lock(m_mutex);
            bool stop_requested = m_cv.wait_for(lock, m_interval, [this]() { return !m_running; });
            
            if (stop_requested) {
                break;
            }

            m_task_queue.push(m_callback);
            lock.unlock();
            m_cv.notify_all();

            processTasks();
        }

        processTasks();
    }

    void processTasks() {
        std::function<void()> task;
        {
            std::lock_guard<std::mutex> lock(m_mutex);
            if (m_task_queue.empty() || m_is_executing) {
                return;
            }
            task = std::move(m_task_queue.front());
            m_task_queue.pop();
            m_is_executing = true;
        }

        if (task) {
            try {
                // 使用共享互斥锁保证全局串行
                std::lock_guard<std::mutex> shared_lock(*m_shared_mutex);
                task();
            } catch (...) {
                std::cerr << "回调函数执行抛出异常\n";
            }
            std::lock_guard<std::mutex> lock(m_mutex);
            m_is_executing = false;
            m_cv.notify_all();
        }
    }

    std::chrono::milliseconds m_interval;
    std::function<void()> m_callback;
    std::queue<std::function<void()>> m_task_queue;
    std::mutex m_mutex;
    std::condition_variable m_cv;
    std::thread m_thread;
    std::atomic<bool> m_running;
    bool m_is_executing; // 受m_mutex保护
    std::shared_ptr<std::mutex> m_shared_mutex; // 全局共享的互斥锁
};

// 跨实例使用示例
int main() {
    auto shared_mutex = std::make_shared<std::mutex>();
    // 两个定时器共享同一个互斥锁,确保myCallback不会并发执行
    CallbackTimer timer1(std::chrono::seconds(2), myCallback, shared_mutex);
    CallbackTimer timer2(std::chrono::seconds(3), myCallback, shared_mutex);
    
    std::this_thread::sleep_for(std::chrono::seconds(20));
    return 0;
}

注意事项

  • 若使用lambda作为回调,由于每个lambda的地址唯一,基于函数地址的自动映射方案会失效,此时必须手动传入共享互斥锁
  • 共享互斥锁的生命周期必须长于所有使用它的定时器实例,避免悬空引用

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 15:27:56