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

vxWorks下如何降低std::condition_variable::wait_for的线程唤醒延迟

问题背景与环境
  • 操作系统:vxWorks 7 23.09
  • C版本:C17
  • 编译器:Clang
  • 硬件平台:英特尔11代i7处理器
  • BSP:英特尔通用64位BSP
问题描述

现有两个线程:以太网读取线程负责从Socket读取数据推入std::queue,随后通知主线程(处理线程)。主线程无数据时通过std::condition_variable::wait_for让出CPU;此前采用忙等待模式。

忙等待时,以太网数据整体处理耗时3-15微秒;改用std::condition_variable::wait_for后,仅线程唤醒阶段就耗时50-250微秒,其余处理步骤仍可在15微秒内完成,忙等待的开销反而低于当前方案。

需要找到降低线程唤醒延迟的策略,可替换std::condition_variable::wait_for,优先选择低延迟/高响应的方案。

测试输出
Calling Notify All : 1731492554671893
Inside Wait For : 1731492554672106
Calling Notify All : 1731492557672289
Inside Wait For : 1731492557672506
Calling Notify All : 1731492560672740
Inside Wait For : 1731492560672901
现有代码实现

wait_for_event.hpp

#include <mutex>
#include <chrono>
#include <iostream>
#include <optional>
#include <condition_variable>

namespace wfe
{
    enum class WaitForEventStatus : std::uint8_t
    {
        Error = 0,
        Timeout,
        NoTimeout,
        NoError  //Required for just wait
    };

    class WaitForEvent
    {
    private:
        std::mutex              mutex;
        std::condition_variable condition_variable;
    public:

        void notify_one(void) noexcept
        {
            this->condition_variable.notify_one();
        }

        void notify_all(void) noexcept
        {
            this->condition_variable.notify_all();
        }

        void wait(void)
        {
            std::unique_lock<std::mutex> lock(this->mutex);
            this->condition_variable.wait(lock);
        }

        template <class Rep, class Period>
        WaitForEventStatus wait_for(const std::chrono::duration<Rep,Period>& wait_time)
        {
            try
            {
                std::unique_lock<std::mutex> lock(this->mutex);
                auto return_status = this->condition_variable.wait_for(lock, wait_time);
                if(return_status == std::cv_status::no_timeout)
                {
                    return WaitForEventStatus::NoTimeout;
                }
                else
                {
                    return WaitForEventStatus::Timeout;
                }
            }
            catch(const std::exception& e)
            {
                std::cerr << e.what() << '\n';
                return WaitForEventStatus::Error;
            }
        }

        template <class Clock, class Duration>
        WaitForEventStatus wait_until(const std::chrono::time_point<Clock,Duration>& abs_time)
        {
            try
            {
                std::unique_lock<std::mutex> lock(this->mutex);
                auto return_status = this->condition_variable.wait_until(lock, abs_time);
                if(return_status == std::cv_status::no_timeout)
                {
                    return WaitForEventStatus::NoTimeout;
                }
                else
                {
                    return WaitForEventStatus::Timeout;
                }
            }
            catch(const std::exception& e)
            {
                std::cerr << e.what() << '\n';
                return WaitForEventStatus::Error;
            }
        }

        WaitForEvent() = default;
        ~WaitForEvent() = default;
    };
}

main.cpp

#include <mutex>
#include <thread>
#include <chrono>
#include <iostream>
#include <condition_variable>
#include "wait_for_event.hpp"

wfe::WaitForEvent event;

void thread_runner(void)
{
    while (true)
    {
        // Code to read data from socket and push the data in queue
        std::this_thread::sleep_for (std::chrono::seconds(3));//Only for testing - not on live code
        std::cout << "Calling Notify All : " << std::chrono::duration_cast<std::chrono::microseconds>
            (std::chrono::system_clock::now().time_since_epoch()).count() << std::endl;
        event.notify_all();
    }
}

int main(int argc, char const *argv[])
{
    std::thread t {thread_runner};
    while (true)
    {
        auto return_status = event.wait_for(std::chrono::seconds(1));
        if(return_status == wfe::WaitForEventStatus::Timeout)
        {
            //std::cout << "WaitForEventStatus::Timeout\n";
        }
        else if(return_status == wfe::WaitForEventStatus::NoTimeout)
        {
            //std::cout << "WaitForEventStatus::NoTimeout\n";
            std::cout << "Inside Wait For : " << std::chrono::duration_cast<std::chrono::microseconds>
            (std::chrono::system_clock::now().time_since_epoch()).count() << std::endl;
        }
    }
    
    return 0;
}
低延迟唤醒策略方案

1. 替换为vxWorks原生同步机制

std::condition_variable是标准库实现,在实时系统中通常会引入额外的内核调度开销。vxWorks提供了更轻量的同步原语,适合低延迟场景:

  • 二进制信号量:使用semBInit()初始化,读取线程用semGive()通知,处理线程用semTake()等待。信号量的内核路径更短,唤醒延迟远低于标准库条件变量。
  • 事件标志:如果需要同时等待多个事件,eventLib的eventReceive()可以实现精准唤醒,且延迟可控。

示例(二进制信号量替换):

#include <semLib.h>

SEM_ID g_dataSem;

// 初始化(主线程启动时)
g_dataSem = semBCreate(SEM_Q_PRIORITY, SEM_EMPTY);

// 读取线程通知
semGive(g_dataSem);

// 处理线程等待
STATUS status = semTake(g_dataSem, WAIT_FOREVER); // 或指定超时
if (status == OK) {
    // 处理数据
}

2. 调整线程优先级与调度策略

在vxWorks中,线程优先级直接影响调度延迟:

  • 将处理线程的优先级设置为高于读取线程,确保处理线程被唤醒后能立即抢占CPU,避免被其他线程阻塞。
  • 采用FIFO调度策略(taskOptionsSet()设置VX_FIFO_SCHED),避免时间片轮转带来的调度延迟。

示例:

// 设置处理线程为FIFO调度,优先级10(数值越小优先级越高)
taskOptionsSet(taskIdSelf(), VX_FIFO_SCHED, TRUE);
taskPrioritySet(taskIdSelf(), 10);

3. 忙等待+短暂休眠的混合模式

完全忙等待会占用CPU,但可以结合极短的休眠平衡CPU占用与延迟:

  • 当队列无数据时,先忙等待若干轮,若仍无数据则调用taskDelay(1)(vxWorks的休眠函数,精度为系统时钟tick)。
  • 这种方式在数据频繁到达时接近忙等待的低延迟,数据稀疏时减少CPU占用。

示例:

void main_thread() {
    while (true) {
        int spin_count = 0;
        bool has_data = false;
        while (spin_count < 1000) { // 忙等待1000轮(根据CPU频率调整)
            std::lock_guard<std::mutex> lock(queue_mutex);
            if (!data_queue.empty()) {
                has_data = true;
                break;
            }
            spin_count++;
        }
        if (has_data) {
            // 处理数据
        } else {
            taskDelay(1); // 休眠1个tick(通常为1ms,若需要更高精度可调整系统时钟)
        }
    }
}

4. 无锁队列+自旋锁

使用无锁队列(如boost::lockfree::queue或自定义CAS实现)替代std::queue+std::mutex,避免锁竞争带来的延迟:

  • 无锁队列通过原子操作实现线程安全,无需内核级锁,唤醒时无需等待锁释放。
  • 处理线程采用自旋等待队列非空,配合内存屏障确保数据可见性。

示例(简化版无锁队列等待):

#include <atomic>
#include <vector>

std::atomic<bool> has_new_data = false;
std::vector<char> data_queue;

void main_thread() {
    while (true) {
        while (!has_new_data.load(std::memory_order_acquire)) {
            // 自旋等待,内存屏障确保读取最新值
        }
        std::lock_guard<std::mutex> lock(queue_mutex); // 若需要批量处理仍需轻量锁
        // 处理队列数据
        has_new_data.store(false, std::memory_order_release);
    }
}

5. 内核级Socket通知

直接利用vxWorks的Socket异步通知机制,避免用户态线程间同步的开销:

  • 使用ioctl()设置Socket的FIOASYNC标志,绑定处理线程的信号量或事件,当Socket有数据到达时,内核直接唤醒处理线程。
  • 这种方式跳过读取线程的中转,直接由内核触发处理,延迟最低。

示例:

#include <socket.h>
#include <ioctl.h>

int sock_fd = socket(AF_INET, SOCK_STREAM, 0);
SEM_ID sock_sem = semBCreate(SEM_Q_PRIORITY, SEM_EMPTY);

// 设置Socket异步通知,关联信号量
ioctl(sock_fd, FIOSETOWN, taskIdSelf());
ioctl(sock_fd, FIOASYNC, TRUE);
ioctl(sock_fd, FIOSEMSET, (char*)&sock_sem);

// 处理线程等待信号量
semTake(sock_sem, WAIT_FOREVER);
// 直接读取Socket数据处理

6. 编译器与系统优化

  • 编译选项:添加-O3优化,开启-march=skylake(针对11代i7),禁用调试符号,减少代码冗余。
  • vxWorks系统配置:关闭不必要的内核特性(如调试、日志),调整系统时钟tick为更小值(如1000Hz即1ms,若硬件支持可更高),减少调度延迟。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 10:47:06