CHERI 能力硬件架构

Linux 内核 cmwq 工作队列深度工程实践:并发管理、优先级反转与生产环境调优全路径解析

工作队列(workqueue)是 Linux 内核中最核心的异步执行机制之一。从块设备的读写完成、网络栈的软中断处理,到容器运行时的 cgroup 通知,几乎所有子系统都依赖 workqueue 来延迟执行工作。然而,当生产环境出现"workqueue stall"或"hung task"时,运维工程师往往束手无策。本文将深入剖析 cmwq(Concurrency Managed Workqueue)的完整架构——从数据结构、并发模型到优先级反转防护,再到生产环境中 hung task 的诊断与调优。

一、为什么我们需要理解 workqueue 内部机制

在一次真实的生产事故中,某金融 Redis 集群在流量突增时出现了随机性 500ms+ 的延迟毛刺。perf top 显示 process_one_work 占用超过 40% 的 CPU,但进一步追踪发现并非 CPU 不足,而是大量 WQ_HIGHPRI 工作项被低优先级任务持有的 mutex 阻塞——典型的优先级反转(Priority Inversion)场景。

这个案例揭示了理解 workqueue 内部机制的必要性:当 workqueue 成为瓶颈时,简单的增加 CPU 核数不仅无法解决问题,反而可能加剧竞争。我们需要从操作系统调度层面理解 workqueue 的工作方式。

二、workqueue 的核心数据结构

2.1 worker_pool:worker 的归属单元

每个 CPU 优先级组合对应一个 worker_pool,这是 worker 的管理单元:

struct worker_pool {
    raw_spinlock_t lock;
    int cpu;                    // 绑定的 CPU(-1 表示未绑定)
    int node;                   // NUMA 节点 ID
    int id;                     // 全局唯一标识

    unsigned int flags;         // WQ_UNBOUND / WQ_HIGHPRI 等
    unsigned int watchdog_touched; // watchdog 心跳计数器

    struct list_head worklist;  // 待处理工作队列
    int nr_workers;            // 当前活跃 worker 数
    int nr_idle;               // 空闲 worker 数

    struct timer_node *idle_timer; // 空闲超时回收定时器
    unsigned long idle_jiffies;    // 空闲超时阈值

    /* 生产环境关键:busy_worker 链表用于 watchdog 检测 */
    struct list_head busy_list;   
    struct hlist_node hash_node;   // 全局 worker_pool 哈希表
};

worker_pool 最重要的设计是多优先级池:同一 CPU 上,WQ_HIGHPRI 和 WQ_NORMAL 各自拥有独立的池。这确保了高优先级 worker 不会被低优先级工作项"饿死"。

2.2 work_struct:工作项的封装

每项工作由一个 work_struct 描述:

struct work_struct {
    atomic_long_t data;         // 状态标志 + owner ID
    struct list_entry entry;    // 链表节点
    work_func_t func;           // 实际执行函数
};

/* data 字段的位分配 */
#define WORK_STRUCT_PENDING     (1UL << 0)  // 已入队但未执行
#define WORK_STRUCT_INACTIVE    (1UL << 1)  // 已执行完毕
#define WORK_STRUCT_COLOR_SHIFT 2           // 颜色位(用于延迟绑定)
#define WORK_STRUCT_COLOR_BITS  4
#define WORK_STRUCT_FLAG_MASK   0xFL        // 低 4 位为标志
#define WORK_STRUCT_WQ_DATA_MASK (~WORK_STRUCT_FLAG_MASK)

关键的 WORK_STRUCT_PENDING 位用于保证幂等性:对同一 work_struct 重复调用 schedule_work() 不会产生重复入队,这是内核中许多子系统依赖的保证。

2.3 worker:实际的执行线程

struct worker {
    struct work_struct *current_work;  // 正在执行的工作
    work_func_t current_func;          // 对应的执行函数
    struct pool_workqueue *pwq;        // 关联的 pwq
    struct list_head node;             // 链入 worker_pool
    struct list_head scheduled;        // 延迟工作的链表

    unsigned int flags;                // WORKER_DIE / WORKER_IDLE 等
    int id;                            // pool 内唯一 ID
    int sleeping;                      // 是否正在睡眠等待任务

    /* affinity 与 NUMA 亲和性 */
    cpumask_t *ptr_mask;              // 最近执行过的 CPU 掩码
    char desc[KSYM_NAME_LEN];          // 线程名称(用于 ps/top)
};

三、cmwq 并发管理的核心算法

3.1 mayday 机制:何时创建新 worker

cmwq 的并发管理核心在于动态创建 worker 的时机判断。当有新工作项需要执行但所有现有 worker 都处于 busy 状态时,系统通过 worker_enter_idle() 和 maybe_create_worker() 协同工作:

static void pool_maybe_idle_timer(struct timer_list *t)
{
    struct worker_pool *pool = from_timer(pool, t, idle_timer);
    bool idle = pool->nr_idle > pool->nr_workers / 2;

    if (idle) {
        /* 空闲过多,标记可回收 */
        mod_timer(&pool->idle_timer, jiffies + IDLE_WORKER_TIMEOUT);
    }
    /* ... */
}

static bool need_more_worker(struct worker_pool *pool)
{
    /* 核心判断逻辑 */
    return !list_empty(&pool->worklist) && !pool->nr_idle;
}

创建新 worker 时,系统执行如下逻辑:

static struct worker *create_worker(struct worker_pool *pool)
{
    struct worker *worker = NULL;
    int id = -1;
    char id_buf[16];

    id = ida_alloc(&pool->worker_ida);
    if (id < 0)
        goto fail;

    worker = alloc_worker(pool->node);
    if (!worker)
        goto fail;

    worker->pool = pool;
    worker->id = id;
    worker->task = kthread_create_on_node(worker_thread, worker,
                                          pool->node,
                                          "kworker/%s-%d",
                                          pool->desc, id);
    if (IS_ERR(worker->task))
        goto fail;

    /* 设置 cgroup 归属和调度策略 */
    set_cpus_allowed_ptr(worker->task, pool->attrs->cpumask);
    set_user_nice(worker->task, pool->attr_nice);

    /* WQ_HIGHPRI 池使用 RT 调度 */
    if (pool->flags & WQ_HIGHPRI) {
        struct sched_param param = { .sched_priority = HIGHPRI_NICE_LEVEL };

        sched_setscheduler_nocheck(worker->task,
                                   SCHED_NORMAL,
                                   &param);
        set_user_nice(worker->task, HIGHPRI_NICE_LEVEL);
    }

    wake_up_process(worker->task);
    return worker;
}

3.2 worker 的忙碌判断与 watchdog

生产环境诊断 workqueue hung 的关键在于 workqueue_watchdog。每个 worker 开始执行工作时会设置 WORKER_PREP 位,完成时清除。watchdog 定期检查所有 pool 中 worker 的活跃状态:

static void wq_watchdog_timer_fn(struct timer_list *unused)
{
    unsigned long threshold = jiffies - WQ_WATCHDOG_MAX_STUCK;
    int cpu;

    /* 遍历所有 worker_pool */
    for_each_online_cpu(cpu) {
        struct worker_pool *pool;

        for_each_pool(pool, cpu) {
            if (pool->flags & WQ_UNBOUND)
                continue;

            struct worker *worker;
            bool found_busy = false;
            
            list_for_each_entry(worker, &pool->busy_list, node) {
                if (time_after(worker->last_active, timestamp))
                    continue;

                /* 找到 stuck worker */
                pr_emerg("BUG: workqueue lockup - pool %d:%s "
                         "appears to be stuck for %lu seconds\n",
                         pool->cpu, pool->node,
                         jiffies_to_msecs(jiffies - worker->last_active) / 1000);
                
                /* dump 调用栈 */
                sched_show_task(worker->task);
                show_stack(worker->task, NULL);
                found_busy = true;
            }
        }
    }

    mod_timer(&wq_watchdog_timer, jiffies + WQ_WATCHDOG_INTERVAL);
}

3.3 unbound workqueue 的 NUMA 亲和性

对于 WQ_UNBOUND 类型的队列,v6.6 内核引入了 NUMA 感知的 worker 分配策略。每个 unbound pool 可以使用 per-NUMA-node 子池,确保 worker 在合适节点的 CPU 上执行:

/* NUMA 节点亲和性选择 */
static void wq_update_unbound_numa(struct workqueue_struct *wq, int cpu)
{
    struct pool_workqueue *pwq;
    int node = cpu_to_node(cpu);

    /* 如果迁移到新 NUMA 节点,选择该节点的 pool */
    rcu_read_lock();
    pwq = rcu_dereference(*wq->cpu_pwq_tbl, node);
    rcu_read_unlock();
    
    if (pwq && pwq->pool->node == node) {
        current->cpus_ptr = pool->attrs->cpumask;
    }
}

四、生产环境实战:优先级反转诊断与修复

4.1 使用 ftrace 追踪 workqueue 延迟

#!/bin/bash
# 追踪 workqueue 入队到执行的时间延迟

TRACE_DIR=/sys/kernel/debug/tracing

# 启用 workqueue 事件
echo 0 > $TRACE_DIR/tracing_on

# 配置追踪点
echo 'workqueue:queue_work' > $TRACE_DIR/set_event
echo 'workqueue:execute_start' > $TRACE_DIR/set_event
echo 'workqueue:execute_end' > $TRACE_DIR/set_event

# 过滤特定 workqueue(可选)
echo 'target_wq="system_highpri"' > $TRACE_DIR/events/workqueue/queue_work/filter

# 开始追踪
echo 1 > $TRACE_DIR/tracing_on

sleep 30

echo 0 > $TRACE_DIR/tracing_on
cat $TRACE_DIR/trace > /tmp/workqueue_trace.txt
cat /tmp/workqueue_trace.txt

4.2 BPF/eBPF 实时 workqueue 监控

// workqueue_latency.bpf.c
#include <linux/bpf.h>
#include <linux/ptrace.h>
#include <linux/workqueue.h>

struct wq_latency_event {
    u64 ts_enqueue;
    u64 ts_execute;
    u32 cpu;
    u32 is_highpri;
    char comm[16];
};

BPF_HASH(start, struct work_struct *, struct wq_latency_event);
BPF_PERF_OUTPUT(events);

TRACEPOINT_PROBE(workqueue, queue_work) {
    struct wq_latency_event ev = {};
    struct work_struct *work = (struct work_struct *)args->work;
    
    ev.ts_enqueue = bpf_ktime_get_ns();
    ev.cpu = bpf_get_smp_processor_id();
    ev.is_highpri = (args->req_cpu == WQ_HIGHPRI_REQ_CPU);
    bpf_get_current_comm(&ev.comm, sizeof(ev.comm));
    
    start.update(&work, &ev);
    return 0;
}

TRACEPOINT_PROBE(workqueue, execute_start) {
    struct work_struct *work = (struct work_struct *)args->work;
    struct wq_latency_event *evp;
    
    evp = start.lookup(&work);
    if (evp) {
        evp->ts_execute = bpf_ktime_get_ns();
        
        u64 latency_ns = evp->ts_execute - evp->ts_enqueue;
        
        /* 只输出延迟超过 100us 的事件 */
        if (latency_ns > 100000) {
            events.perf_submit(ctx, evp, sizeof(*evp));
        }
    }
    return 0;
}

4.3 优先级反转修复方案

当发现 WQ_HIGHPRI 工作项被低优先级任务阻塞时,有三种修复策略:

方案一:使用 per-CPU BH 替代 workqueue

对于在中断下半部执行的轻量级工作,使用 local_bh_disable()/enable() + 私有 BH 链表,完全规避 workqueue 延迟:

static DEFINE_PER_CPU(struct my_bh_work, my_bh_queue);

static void my_bh_handler(struct irqlet_struct *bh)
{
    struct my_bh_work *work = this_cpu_ptr(&my_bh_queue);
    
    while (!list_empty(&work->head)) {
        struct my_work *w = list_first_entry(&work->head, struct my_work, list);
        list_del(&w->list);
        w->func(w->data);
    }
}

/* 在硬中断中调用 */
static irqreturn_t my_irq_handler(int irq, void *dev)
{
    struct my_work *w = get_work_struct();
    list_add_tail(&w->list, &this_cpu_ptr(&my_bh_queue)->head);
    irq_or_softirq_bh_raise(my_bh_bit);
    return IRQ_HANDLED;
}

方案二:创建隔离的 WQ_UNBOUND_HIGHPRI 队列

创建专用的 unbound highpri workqueue,绑定到指定 CPU 组:

#include <linux/workqueue.h>
#include <linux/sched.h>

static struct workqueue_struct *wq_isolate;
static struct work_struct my_work;

static void my_work_func(struct work_struct *work)
{
    /* 关键路径执行 */
}

static void setup_isolated_wq(void)
{
    struct workqueue_attrs attrs;
    
    /* 创建 unbound workqueue */
    wq_isolate = alloc_workqueue("my_wq_isolated", 
                                  WQ_UNBOUND | WQ_HIGHPRI, 0);
    
    /* 配置 CPU 隔离:只使用 CPU 4-7 */
    workqueue_attrs_init(&attrs);
    attrs.nice = HIGHPRI_NICE_LEVEL;
    cpumask_clear(attrs.cpumask);
    cpumask_set_cpu(4, attrs.cpumask);
    cpumask_set_cpu(5, attrs.cpumask);
    cpumask_set_cpu(6, attrs.cpumask);
    cpumask_set_cpu(7, attrs.cpumask);
    
    workqueue_set_unbound_cpumask(wq_isolate, &attrs);
}

方案三:使用 io_uring 替代 workqueue 进行延迟任务

v6.7+ 内核引入的 io_uring 链接操作可以部分替代 workqueue 的功能,特别是在 I/O 场景:

/* 使用 io_uring 的链接任务替代 delayed_work */
struct io_uring_sqe *sqe, *next_sqe;

sqe = io_uring_get_sqe(ring);
io_uring_prep_rw(IORING_OP_READ, sqe, fd, buf, len, offset);
sqe->flags |= IOSQE_IO_LINK;  // 设置链接标志

next_sqe = io_uring_get_sqe(ring);
io_uring_prep_rw(IORING_OP_WRITE, next_sqe, out_fd, buf, len, out_offset);
next_sqe->user_data = (u64)completion_callback; // 完成回调

io_uring_submit(ring);

五、常见问题诊断速查表

5.1 workqueue hung 排查流程

1. 确认症状
   └─ dmesg | grep "workqueue lockup"
   └─ 检查 watchdog: cat /sys/kernel/debug/workqueue/state

2. 定位 stuck worker  
   └─ ps -eo pid,stat,comm | grep kworker
   └─ cat /proc/<pid>/wchan  
   └─ cat /proc/<pid>/stack  

3. 分析阻塞链
   └─ 检查 lockdep: cat /proc/lockdep_chains
   └─ 检查 mutex/spinlock 持有者

4. 确定根因
   └─ 是否为优先级反转?
   └─ 是否为不可中断睡眠?
   └─ 是否为上报时间戳异常?

5. 实施修复
   └─ 调整 nice 值 / RT priority
   └─ 隔离 CPU 资源
   └─ 重构为 per-CBH 或 io_uring

5.2 sysctl 关键参数调优

# 增加 workqueue watchdog 超时阈值(默认 30s)
sysctl kernel.workqueue_watchdog_thresh=60000

# 增加单个 worker 允许的工作项数
# (防止频繁创建/销毁 worker 导致 CPU 抖动)
sysctl kernel.workqueue_max_active=512

# unbound workqueue 的 CPU 隔离范围
# 注意:需要重启后生效
echo "4-7" > /sys/devices/virtual/workqueue/cpumask

# 查看当前 workqueue 状态
cat /sys/devices/virtual/workqueue/*/cpumask
cat /sys/devices/virtual/workqueue/*/nice
cat /sys/devices/virtual/workqueue/*/max_active

六、内核演进趋势

从 v5.10 到 v6.8 内核,workqueue 的演进方向主要集中在三个维度:

  1. 更智能的并发管理:v6.3 引入的 workqueue.affinity_strictness 允许管理员在 NUMA 亲和性和负载均衡之间做权衡。
    1. 内存效率提升:v6.5 开始对空闲的 worker_pool 执行延迟回收,减少每个池的常驻内存开销(尤其在大型 NUMA 系统上)。
      1. 与 cgroup 深度集成:v6.7 引入 cgroup 级别的 workqueue 带宽限制,防止单个 cgroup 占用过多 worker 资源。
      2. # 查看 cgroup v2 的 workqueue 限制(v6.7+)
        cat /sys/fs/cgroup/<cgroup>/workqueue.max_active
        cat /sys/fs/cgroup/<cgroup>/workqueue.nice

        七、总结

        工作队列作为内核异步执行的基石,其性能直接影响整个系统的响应延迟与吞吐量。本文从数据结构、并发管理算法、优先级反转防护、hung task 诊断四个维度进行了系统性分析。核心要点:

        • cmwq 是动态的:worker 根据负载自动创建/销毁,不可简单等同于固定线程池
        • 优先级隔离是必须的:WQ_HIGHPRI 池与 normal 池的分离阻塞了"低优先级工作饿死高优先级工作"的路径
        • unbound 队列适合 NUMA 密集型工作:通过 per-node 亲和性避免跨 NUMA 调度开销
        • 诊断链路:ftrace → perf → BPF 三重工具覆盖从宏观趋势到单个工作项的全景观测

        对于运维和性能工程师来说,理解 workqueue 不仅有助于解决当前的性能问题,更是理解 Linux 内核"异步执行"哲学的一把钥匙——从 softirq 到 tasklet,从 threaded irq 到 io_uring,内核的异步演进始终围绕同一个核心矛盾:在保证数据一致性的前提下,最小化同步等待的开销。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部