执行队列

qianmoQqianmoQ· 更新于 2026-10-05· 阅读 16 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

执行队列

高性能执行队列。

概述

类似于 kylin 的 ExecMan,ExecutionQueue]提供了异步串行执行的功能。ExecutionQueue 的相关技术最早使用在 RPC 中实现多线程向同一个 fd 写数据]。在 r31345 之后加入到 bthread。ExecutionQueue 提供了如下基本功能:

  • 异步有序执行:任务在另外一个单独的线程中执行,并且执行顺序严格和提交顺序一致。
  • 多生产者(Multi Producer):多个线程可以同时向同一个 ExecutionQueue 提交任务
  • 支持 cancel 一个已经提交的任务
  • 支持 stop
  • 支持高优任务插队

和 ExecMan 的主要区别:

  • ExecutionQueue 的任务提交接口是wait-free]的,ExecMan 依赖了 lock,这意味着当机器整体比较繁忙的时候,使用 ExecutionQueue 不会因为某个进程被系统强制切换导致所有线程都被阻塞。
  • ExecutionQueue 支持批量处理:执行线程可以批量处理提交的任务,获得更好的 locality。ExecMan 的某个线程处理完某个 AsyncClient 的 AsyncContext 之后,下一个任务很可能是属于另外一个 AsyncClient 的 AsyncContext,这时候 CPU cache 会在不同 AsyncClient 依赖的资源间进行不停的切换。
  • ExecutionQueue 的处理函数不会被绑定到固定的线程中执行,ExecMan 中是根据 AsyncClient hash 到固定的执行线程。不同的 ExecutionQueue 之间的任务处理完全独立,当线程数足够多的情况下,所有非空闲的 ExecutionQueue 都能同时得到调度。同时也意味着当线程数不足的时候,ExecutionQueue 无法保证公平性,当发生这种情况的时候需要动态增加 bthread 的 worker 线程来增加整体的处理能力。
  • ExecutionQueue 运行线程为 bthread,可以随意地使用一些 bthread 同步原语而不用担心阻塞 pthread 的执行;而在 ExecMan 里面得尽量避免使用较高概率会导致阻塞的同步原语。

背景

在多核并发编程领域,Message passing]作为一种解决竞争的手段得到了比较广泛的应用。它按照业务依赖的资源将逻辑拆分成若干个独立的 actor,每个 actor 负责对应资源的维护工作。当一个流程需要修改某个资源的时候,就转化为一个消息发送给对应 actor,这个 actor(通常在另外的上下文中)根据命令内容对这个资源进行相应的修改,之后可以选择唤醒调用者(同步)或者提交到下一个 actor(异步)的方式进行后续处理。

img

ExecutionQueue 与 Mutex 对比

ExecutionQueue 和 mutex 都可以用来在多线程场景中消除竞争。相比较使用 mutex,使用 ExecutionQueue 有着如下几个优点:

  • 角色划分比较清晰,概念理解上比较简单,实现中无需考虑锁带来的问题(比如死锁)
  • 能保证任务的执行顺序,mutex 的唤醒顺序不能得到严格保证。
  • 所有线程各司其职,都能在做有用的事情,不存在等待。
  • 在繁忙、卡顿的情况下能更好的批量执行,整体上获得较高的吞吐。

但是缺点也同样明显:

  • 一个流程的代码往往散落在多个地方,代码理解和维护成本高。
  • 为了提高并发度,一件事情往往会被拆分到多个 ExecutionQueue 进行流水线处理,这样会导致在多核之间不停地进行切换,会付出额外的调度以及同步 cache 的开销,尤其是竞争的临界区非常小的情况下,这些开销不能忽略。
  • 同时原子地操作多个资源实现会变得复杂,使用 mutex 可以同时锁住多个 mutex,用了 ExecutionQueue 就需要依赖额外的 dispatch queue 了。
  • 由于所有操作都是单线程的,某个任务运行慢了就会阻塞同一个 ExecutionQueue 的其他操作。
  • 并发控制变得复杂,ExecutionQueue 可能会由于缓存的任务过多占用过多的内存。

不考虑性能和复杂度,理论上任何系统都可以只使用 mutex 或者 ExecutionQueue 来消除竞争。但是复杂系统的设计上,建议根据不同的场景灵活决定如何使用这两个工具:

  • 如果临界区非常小,竞争又不是很激烈,优先选择使用 mutex,之后可以结合contention profiler]来判断 mutex 是否成为瓶颈。
  • 需要有序执行,或者存在无法消除的激烈竞争但可以通过批量执行来提高吞吐,可以选择使用 ExecutionQueue。

总之,多线程编程没有万能的模型,需要根据具体的场景,结合丰富的 profiling 工具,最终在复杂度和性能之间找到合适的平衡。

特别指出一点,Linux 中 mutex 无竞争的 lock/unlock 只需要几条原子指令,在绝大多数场景下的开销都可以忽略不计。

使用方式

实现执行函数

// Iterate over the given tasks
//
// Example:
//
// #include <bthread/execution_queue.h>
//
// int demo_execute(void* meta, TaskIterator<T>& iter) {
//     if (iter.is_stopped()) {
//         // destroy meta and related resources
//         return 0;
//     }
//     for (; iter; ++iter) {
//         // do_something(meta, *iter)
//         // or do_something(meta, iter->a_member_of_T)
//     }
//     return 0;
// }
template <typename T>
class TaskIterator;

启动一个 ExecutionQueue:

// Start a ExecutionQueue. If |options| is NULL, the queue will be created with
// default options.
// Returns 0 on success, errno otherwise
// NOTE: type |T| can be non-POD but must be copy-constructible
template <typename T>
int execution_queue_start(
        ExecutionQueueId<T>* id,
        const ExecutionQueueOptions* options,
        int (*execute)(void* meta, TaskIterator<T>& iter),
        void* meta);

创建的返回值是一个 64 位的 id,相当于 ExecutionQueue 实例的一个弱引用](https://en.wikipedia.org/wiki/Weak_reference),可以 wait-free 地在 O(1) 时间内定位一个 ExecutionQueue。你可以到处拷贝这个 id,甚至可以把它放进 RPC 中,作为远端资源的定位工具。你必须保证 meta 的生命周期,在对应的 ExecutionQueue 真正停止前不会被释放。

停止一个 ExecutionQueue:

// Stop the ExecutionQueue.
// After this function is called:
//  - All the following calls to execution_queue_execute would fail immediately.
//  - The executor will call |execute| with TaskIterator::is_queue_stopped() being
//    true exactly once when all the pending tasks have been executed, and after
//    this point it's ok to release the resource referenced by |meta|.
// Returns 0 on success, errno othrwise
template <typename T>
int execution_queue_stop(ExecutionQueueId<T> id);

// Wait until the the stop task (Iterator::is_queue_stopped() returns true) has
// been executed
template <typename T>
int execution_queue_join(ExecutionQueueId<T> id);

stop和join都可以多次调用,都会有合理的行为。stop可以随时调用,而不必担心线程安全性问题。

和fd的close类似,如果stop不被调用,相应的资源会永久泄露。

释放meta的安全时机:可以在execute函数中收到 iter.is_queue_stopped()==true 的任务时释放,也可以等到join返回之后释放。注意不要double-free。

提交任务

struct TaskOptions {
    TaskOptions();
    TaskOptions(bool high_priority, bool in_place_if_possible);

    // Executor would execute high-priority tasks in the FIFO order but before
    // all pending normal-priority tasks.
    // NOTE: We don't guarantee any kind of real-time as there might be tasks still
    // in process which are uninterruptible.
    //
    // Default: false
    bool high_priority;

    // If |in_place_if_possible| is true, execution_queue_execute would call
    // execute immediately instead of starting a bthread if possible
    //
    // Note: Running callbacks in place might cause the dead lock issue, you
    // should be very careful turning this flag on.
    //
    // Default: false
    bool in_place_if_possible;
};

const static TaskOptions TASK_OPTIONS_NORMAL = TaskOptions(/*high_priority=*/ false, /*in_place_if_possible=*/ false);
const static TaskOptions TASK_OPTIONS_URGENT = TaskOptions(/*high_priority=*/ true, /*in_place_if_possible=*/ false);
const static TaskOptions TASK_OPTIONS_INPLACE = TaskOptions(/*high_priority=*/ false, /*in_place_if_possible=*/ true);

// Thread-safe and Wait-free.
// Execute a task with defaut TaskOptions (normal task);
template <typename T>
int execution_queue_execute(ExecutionQueueId<T> id,
                            typename butil::add_const_reference<T>::type task);

// Thread-safe and Wait-free.
// Execute a task with options. e.g
// bthread::execution_queue_execute(queue, task, &bthread::TASK_OPTIONS_URGENT)
// If |options| is NULL, we will use default options (normal task)
// If |handle| is not NULL, we will assign it with the hanlder of this task.
template <typename T>
int execution_queue_execute(ExecutionQueueId<T> id,
                            typename butil::add_const_reference<T>::type task,
                            const TaskOptions* options);
template <typename T>
int execution_queue_execute(ExecutionQueueId<T> id,
                            typename butil::add_const_reference<T>::type task,
                            const TaskOptions* options,
                            TaskHandle* handle);

high_priority 之间的任务执行顺序同样严格按照提交顺序,这一点与 ExecMan 不同:ExecMan 的 QueueExecEmergent 中 AsyncContext 的执行顺序是未定义的(undefined)。但这也意味着,你无法将任何任务插队到一个高优先级任务之前执行。

开启 inplace_if_possible,在无竞争的场景中可以省去一次线程调度和 cache 同步的开销。但可能会造成死锁或者递归层数过多(比如不停地 ping-pong)等问题,开启前请先确认你的代码中不存在这些问题。

取消一个已提交任务

/// [Thread safe and ABA free] Cancel the corresponding task.
// Returns:
//  -1: The task was executed or h is an invalid handle
//  0: Success
//  1: The task is executing
int execution_queue_cancel(const TaskHandle& h);

返回非0仅仅意味着ExecutionQueue已经将对应的task递给过execute, 真实的逻辑中可能将这个task缓存在另外的容器中,所以这并不意味着逻辑上的task已经结束,你需要在自己的业务上保证这一点.


最后修改于 2023 年 6 月 28 日:typo (devlive-community/knowforge#151) (49251654f)

评论

登录后参与评论

正在加载评论…