Workqueue 机制解析

作者:CherryYang 发布时间: 2026-07-14 阅读量:4 评论数:0

嵌入式存储系统中的 Workqueue 机制解析——以 dpax_sched_work 为例

核心结论dpax_sched_init_work / dpax_sched_add_work 实现了一套仿照 Linux kernel workqueue 设计的单线程延迟执行机制,核心思想是"线程只创建一次,任务无限复用"。
讲人话:不想每次干活都临时招人(创建线程),干脆养一个常驻员工(worker 线程),有活就往他桌上放(链表入队),拍一下他(信号量唤醒)就行。

问题背景:为什么不直接创建线程

假设你正在处理一个请求,中途需要异步发送一次告警事件:

void handle_something() {
    do_main_work();
    send_alarm_event();  // 如果同步执行,当前线程会被阻塞
}

最直觉的做法是 pthread_create 起一个新线程去跑 send_alarm_event,但每次都创建线程的代价很重——栈分配、调度器注册、线程销毁,对于一个只需要跑几十行代码的短任务来说完全不值得。

更好的做法是:系统启动时预先创建好一个 worker 线程,有任务就往队列里扔一个节点,worker 线程自己取出来执行。提交任务的开销只有一次链表插入和一次信号量 up,没有任何线程创建开销——这就是 Workqueue(工作队列)模式。

核心数据结构:任务卡片 dpax_sched_work_s

每个异步任务在提交前,都需要先初始化一个 dpax_sched_work_s 结构体,可以理解为一张"任务卡片":

void dpax_sched_init_work(dpax_sched_work_s *v_pstWork, void (*func)(void *), void *data)
{
    s32 pid;
    pid = dpax_thrd_getpid();
    v_pstWork->siPid = pid;
    v_pstWork->uiTimeout = 0;
    DPAX_INIT_LIST_NODE(&v_pstWork->stNode);
    v_pstWork->pData = data;
    v_pstWork->pfnWorkHandler = func;
}

各字段的作用:

  • siPid:记录是哪个线程/进程提交的任务,用于调试和超时告警时定位问题。
  • uiTimeout:超时时间,为 0 表示不做超时监控。
  • stNode:侵入式链表节点(Intrusive List Node),用于挂到全局任务队列上。DPAX_INIT_LIST_NODE 会把 next/prev 设为特殊的"毒值"(类似 Linux 的 LIST_POISON1),标记此节点当前不在任何链表中。
  • pDatavoid * 万能指针,存储任务执行时需要的参数,回调内部自行转换为具体类型。
  • pfnWorkHandler:函数指针,指向任务真正要执行的函数。

函数指针 void (*func)(void *) 是怎么工作的?

这是 C 语言实现"延迟调用"的标准手法。函数编译后是 .text 段中的一段机器码,函数指针存的就是这段代码的起始地址:

// 注册阶段:调用方填入函数地址和参数
v_pstWork->pfnWorkHandler = send_timeout_alarm_event;  // 存函数地址
v_pstWork->pData          = (void*)work;               // 存参数地址

// 执行阶段:worker 线程通过指针调用
v_pstWork->pfnWorkHandler(v_pstWork->pData);
// 等价于:send_timeout_alarm_event((void*)work);

worker 线程不需要知道具体执行什么函数,只管通过指针调用——这就是 Workqueue 框架与具体业务逻辑解耦的关键。

void * 参数使得框架不必为每种任务定义不同的函数签名,任何类型的数据都可以通过 void * 传入,在回调内部强制转换回具体类型。这是 C 语言泛型编程的标准做法。

任务提交:dpax_sched_add_work 逐段拆解

void dpax_sched_add_work(dpax_sched_work_s *v_pstWork)
{
    if (NULL == v_pstWork) {
        DP_OSAX_LOG(DPLOG_LVL_ERROR, "The v_pstWork is equal NULL.");
        return;
    }

    /* 工作队列自动初始化 */
    if (NOT_INITED == m_stDpaxSchedInitLock.initedOnce ||
        NOT_INITED == m_stDpaxSchedInitLock.initedMulti) {
        if (RET_OK != sched_work_init()) {
            DP_OSAX_LOG(DPLOG_LVL_ERROR, "Adding work failed, as the result of initing sched_work failed");
            return;
        }
    }

    /* 进程处于卸载流程,阻止上层在卸载流程中添加任务 */
    dpax_spinlock_lock(&(m_stDpaxDestroyInfo.destroyLock));
    if (DPAX_SCHED_WORK_IN_DESTROYING == m_stDpaxDestroyInfo.isInDestroying) {
        DP_OSAX_LOG(DPLOG_LVL_ERROR, "In destroying, adding v_pstWork failed.");
    }
    dpax_spinlock_unlock(&(m_stDpaxDestroyInfo.destroyLock));

    if (!IS_LIST_NODE_INIT(&v_pstWork->stNode)) {
        DP_OSAX_LOG(DPLOG_LVL_ERROR, "Work in list is not initial, PID = (%d)", v_pstWork->siPid);
        return;
    }

    dpax_spinlock_lock(&(m_stSchedListInfo.schedListLock));
    LIST_ADD_PREV_CHECK(&(m_stSchedListInfo.schedList), DpaxSchedWorkShow, LIST_ABNORMAL_RET);
    dpax_list_add_tail(&v_pstWork->stNode, &(m_stSchedListInfo.schedList));
    dpax_spinlock_unlock(&(m_stSchedListInfo.schedListLock));
    dpax_sem_enhance_up(&m_stDpaxSchedSema);

    return;
}

整个函数做了四件事:

1. 懒初始化(Lazy Initialization)

if (NOT_INITED == m_stDpaxSchedInitLock.initedOnce ||
    NOT_INITED == m_stDpaxSchedInitLock.initedMulti) {
    sched_work_init();
}

Workqueue 的全局状态(链表、锁、信号量、worker 线程)不是在 main() 里统一初始化的,而是在第一次提交任务时按需创建。initedOnceinitedMulti 两个标志位区分了单进程模式和多进程模式的初始化状态。

2. 卸载保护

dpax_spinlock_lock(&(m_stDpaxDestroyInfo.destroyLock));
if (DPAX_SCHED_WORK_IN_DESTROYING == m_stDpaxDestroyInfo.isInDestroying) {
    DP_OSAX_LOG(DPLOG_LVL_ERROR, "In destroying, adding v_pstWork failed.");
}
dpax_spinlock_unlock(&(m_stDpaxDestroyInfo.destroyLock));

当进程正在执行卸载流程时,worker 线程可能已经退出或正在退出,此时继续往队列里塞任务可能导致任务永远无法执行或访问已释放的资源。通过检查 isInDestroying 标志位,在卸载阶段拦截新任务。

3. 防重复入队

if (!IS_LIST_NODE_INIT(&v_pstWork->stNode)) {
    DP_OSAX_LOG(DPLOG_LVL_ERROR, "Work in list is not initial, PID = (%d)", v_pstWork->siPid);
    return;
}

检查 stNode 是否处于"毒值"状态(即 DPAX_INIT_LIST_NODE 初始化后的状态)。如果不是,说明这个 work 已经在链表里了,重复入队会破坏链表结构。这是对 init_work 中毒值初始化的呼应。

4. 入队 + 唤醒

dpax_spinlock_lock(&(m_stSchedListInfo.schedListLock));
LIST_ADD_PREV_CHECK(...);  // 链表完整性校验
dpax_list_add_tail(&v_pstWork->stNode, &(m_stSchedListInfo.schedList));
dpax_spinlock_unlock(&(m_stSchedListInfo.schedListLock));
dpax_sem_enhance_up(&m_stDpaxSchedSema);

加锁保护全局链表,将 work 节点挂到链表尾部,解锁后执行 dpax_sem_enhance_up 将信号量计数 +1,唤醒阻塞在 sem_wait 上的 worker 线程。

Worker 线程:系统的另一半

代码中没有直接展示 worker 线程的主循环,但从 sched_work_init 创建线程 + sched_work_run_task 执行任务的逻辑可以还原出它的核心结构:

// worker 线程主循环(伪代码,基于实际代码推断)
static void *worker_thread_func(void *arg)
{
    while (1) {
        // 没有任务时,睡在这里,不占用 CPU
        dpax_sem_enhance_down(&m_stDpaxSchedSema);

        // 被唤醒后,从链表取出一个任务
        dpax_spinlock_lock(&(m_stSchedListInfo.schedListLock));
        work = list_first_entry(&(m_stSchedListInfo.schedList), ...);
        dpax_list_del_init(&work->stNode);  // 从链表摘除,stNode 重置为毒值
        dpax_spinlock_unlock(&(m_stSchedListInfo.schedListLock));

        // 执行任务(带超时监控)
        sched_work_run_task(work, index);
    }
}

sched_work_run_task 的实际代码如下:

void sched_work_run_task(dpax_sched_work_s* v_pstWork, uint32_t index)
{
    uint32_t timeout;
    if (NULL == v_pstWork->pfnWorkHandler) {
        DP_OSAX_LOG(DPLOG_LVL_ERROR, "NULL sched work data handler, PID = (%d).", v_pstWork->siPid);
        return;
    }
    timeout = v_pstWork->uiTimeout;
    if (timeout != 0) {
        dpax_update_sched_work_timeout_info(index, dpax_get_jiffies(), timeout * HZ, v_pstWork->siPid);
    }
    v_pstWork->pfnWorkHandler(v_pstWork);
    if (timeout != 0) {
        dpax_update_sched_work_timeout_info(index, 0, 0, 0);
    }
}

如果 uiTimeout 不为 0,执行前后会更新超时监控信息,用于检测任务是否卡死。执行本身就是通过函数指针调用:v_pstWork->pfnWorkHandler(v_pstWork)

信号量的角色:不是阻塞调用方,而是唤醒 Worker

一个容易产生的误解是:dpax_sem_enhance_up(底层对应 sem_post)会不会阻塞提交任务的调用方?

答案是不会。信号量的两个操作是不对称的:

  • sem_post(sema_up):永不阻塞。只是将计数原子 +1,如果有线程在等待则唤醒它。
  • sem_wait(sema_down):会阻塞。计数为 0 时睡眠等待,直到有人调用 sem_post

信号量的计数器还起到了"任务缓冲"的作用——即使 worker 线程还没来得及处理,连续提交多个任务也不会丢失:

初始计数 = 0

提交任务1 → sem_post → 计数变1
提交任务2 → sem_post → 计数变2(worker 还没醒来处理)
提交任务3 → sem_post → 计数变3

worker 醒来:
  sem_wait → 计数变2,取任务1执行
  sem_wait → 计数变1,取任务2执行
  sem_wait → 计数变0,取任务3执行
  sem_wait → 计数=0,阻塞,等下一个任务

延伸:如果是 LWT 协程环境呢?

这套代码的 worker 是真实的 OS 线程(从 dpax_thrd_getpid()spin_lock 等 OS 级原语可以确认),所以 sem_wait 只阻塞这一个 OS 线程,没有问题。

但如果是 M:N 协程模型(多个轻量级线程复用少数 OS 线程),直接调用 sem_wait 会阻塞整个 OS 线程,导致同一 OS 线程上的其他协程全部卡死。这也是为什么真正的协程运行时(如 Go runtime)必须把所有阻塞系统调用替换为协程级别的挂起操作。

实际上,send_alarm_event 的注释也明确写了这一点:

// 起一个非 lwt 异步线程发送关键事件

正是因为告警事件的发送涉及 dpax_system_s 等可能阻塞的操作,不能在 LWT 协程中执行,所以要通过 Workqueue 投递到一个真实的 OS 线程上去跑。

实际使用场景:异步告警事件发送

以 CCDB 超时告警为例,看这套机制如何被业务代码使用:

static void send_alarm_event()
{
    // 限频上报,一天一次
    uint32_t one_day = 24 * 60 * 60;
    if (CCDB_IS_LIMIT_ONCE_PER_INTERVAL(one_day) == TRUE) {
        ccdb_info_log_limit("Send timeout alarm event is limited.");
        return;
    }

    // 堆上分配 work(不能用栈变量,因为函数返回后栈帧销毁)
    dpax_sched_work_s *work = (dpax_sched_work_s*)dpax_malloc((u32)sizeof(dpax_sched_work_s));
    if (work != NULL) {
        LVOS_INIT_WORK(work, send_timeout_alarm_event_callback, (void*)work);
        LVOS_SchedWork(work);
    }
}
void send_timeout_alarm_event(void *p)
{
    // 进来即释放
    dpax_free(p);

    int32_t ret = RETURN_OK;
    const int32_t MAX_STRING_SIZE = 1024;
    char key_event_send_string[MAX_STRING_SIZE];
    ret = snprintf_s(key_event_send_string, sizeof(key_event_send_string),
                     sizeof(key_event_send_string),
                     "%s commonEvent 19910200013 14 \"ccdb access failed.\"",
                     CCDB_ALARM_SCRIPT_PATH);
    CHECK_SNPRINTF_S_RETURN(ret);
    ret = dpax_system_s(key_event_send_string, NULL);
    if (ret == RETURN_OK) {
        ccdb_info_log_limit("Ccdb timeout alarm event send success.");
    } else {
        ccdb_warn_log_limit("Ccdb timeout alarm event send failed, ret:(%d).", ret);
    }
}

这里有一个值得注意的模式:LVOS_INIT_WORK(work, callback, (void*)work)work 自身作为 data 传入回调。回调收到的 void *p 就是 work 指针本身,目的是让回调负责释放这块堆内存。

为什么在回调开头释放而非结尾? 因为后续逻辑中 CHECK_SNPRINTF_S_RETURN 等宏可能触发提前 return,如果 dpax_free 放在结尾,提前返回就会导致内存泄漏。放在开头保证无论后续走哪条路径,内存一定被释放。

如果任务需要传递额外数据,标准做法是定义一个包含 dpax_sched_work_s 的更大结构体:

typedef struct {
    dpax_sched_work_s stWork;  // 放第一个成员
    int               extraData;
    char              msg[64];
} my_work_s;

my_work_s *w = dpax_malloc(sizeof(my_work_s));
w->extraData = 42;
LVOS_INIT_WORK(&w->stWork, my_callback, w);
LVOS_SchedWork(&w->stWork);

void my_callback(void *p) {
    my_work_s *w = (my_work_s *)p;
    use(w->extraData);
    dpax_free(w);
}

术语对应:与 Linux Kernel Workqueue 的关系

这套代码的命名和设计直接对标 Linux kernel workqueue:

Linux Kernel dpax 实现 作用
INIT_WORK() dpax_sched_init_work() 初始化 work item
schedule_work() dpax_sched_add_work() 提交 work item 到队列
struct work_struct dpax_sched_work_s 任务描述结构体
workqueue_struct m_stSchedListInfo(链表+锁)+ m_stDpaxSchedSema(信号量) 队列全局状态
worker thread sched_work_init() 中创建的线程 消费者线程

从设计模式角度,这套机制涉及的术语包括:

  • Producer-Consumer Pattern(生产者消费者模式):调用方生产任务,worker 线程消费。
  • Thread Pool(线程池):此处是退化为单线程的线程池。
  • Deferred Execution(延迟执行):任务不在调用方线程中立即执行,而是延迟到 worker 线程异步执行。

总结

普通写法:每次创建线程

  1. 有异步任务时调用 pthread_create
  2. 执行完毕后线程销毁
  3. 下次又要 pthread_create,栈分配、调度器注册、销毁回收……
  4. 结果:高频场景下线程创建/销毁的开销不可接受

Workqueue 写法:预创建线程 + 任务队列

  1. 系统启动时(或首次提交时懒初始化)创建好 worker 线程
  2. 有任务时只做一次链表尾插 + 信号量 up
  3. worker 线程自动取出执行,执行完继续等待
  4. 结果:提交路径极短(无内存分配、无线程创建),worker 线程永久复用

Workqueue 的本质:把"创建执行者"的开销分摊到系统初始化阶段,运行时只传递"做什么"(函数指针)和"用什么"(void * 参数),实现接近零开销的异步任务分发。

评论