上次说到 GCD 中我们可以通过 GCD 提供的函数创建 block 这种方式,实现对 block 的管理,比如取消 block,比如监听 block 的执行情况,而 block 的 wait/notify 的实现就依赖于 dispatch group 这一功能,所以现在来看下 dispatch group。
功能
总的来说,dispatch group 的存在,就是为了完成对 block 的 wait/notify 这样一个功能,我们可以通过 dispatch_group_async 提交一个 block,然后再通过 dispatch_group_notify 添加一个对 group 的监听,当与 group 关联的所有任务执行完了之后,这个用于监听的 block 就会执行,或者也可以直接使用 dispatch_group_wait 阻塞当前线程等待 group 中任务的执行。
除此之外,group 还提供了 dispatch_group_enter/dispatch_group_leave 这样一对函数,分别用于标记一次 block 执行的开始和结束,dispatch_group_async 的实现其实就是在 block 执行前后自动调用了这两个函数,然后当 group 关联的所有任务都执行完了之后(也就是 enter/leave 函数成对调用之后),所有添加到 group 中的监听 block,就会依次被执行。
下面就看下具体的实现。
实现
创建
dispatch_group_t
dispatch_group_create(void)
{
return _dispatch_group_create_with_count(0);
}
static inline dispatch_group_t
_dispatch_group_create_with_count(uint32_t n)
{
dispatch_group_t dg = _dispatch_object_alloc(DISPATCH_VTABLE(group),
sizeof(struct dispatch_group_s));
dg->do_next = DISPATCH_OBJECT_LISTLESS;
dg->do_targetq = _dispatch_get_default_queue(false);
if (n) {
os_atomic_store2o(dg, dg_bits,
(uint32_t)-n * DISPATCH_GROUP_VALUE_INTERVAL, relaxed);
os_atomic_store2o(dg, do_ref_cnt, 1, relaxed); // <rdar://22318411>
}
return dg;
}
很简单,就是创建了一个实例,分配了下内存。这里的 n 就是指当前的 enter 数量,默认是 0,如果 n 不为 0 的话,就相当于执行了 n 次的 dispatch_group_enter。
提交任务
当我们使用 dispatch_group_async 将一个任务与 group 绑定的时候,本质上就是在 block 提交给 dispatch queue 之前,执行一下 dispatch_group_enter,然后在 block 执行完了之后,执行一下 dispatch_group_leave,具体的可以从代码实现一窥:
dispatch_group_async(dispatch_group_t dg, dispatch_queue_t dq,
dispatch_block_t db)
{
dispatch_continuation_t dc = _dispatch_continuation_alloc();
uintptr_t dc_flags = DC_FLAG_CONSUME | DC_FLAG_GROUP_ASYNC;
dispatch_qos_t qos;
qos = _dispatch_continuation_init(dc, dq, db, 0, dc_flags);
_dispatch_continuation_group_async(dg, dq, dc, qos);
}
static inline void
_dispatch_continuation_group_async(dispatch_group_t dg, dispatch_queue_t dq,
dispatch_continuation_t dc, dispatch_qos_t qos)
{
dispatch_group_enter(dg);
dc->dc_data = dg;
_dispatch_continuation_async(dq, dc, qos, dc->dc_flags);
}
它的实现方式与 dispatch_async 其实基本一致,不同的就是在 dispatch_group_async 中,多加了一个 flag DC_FLAG_GROUP_ASYNC,然后在 _dispatch_continuation_group_async 中执行了下 dispatch_group_enter,并将 group 保存下来,而 leave 的执行可以从 _dispatch_continuation_invoke_inline 中找到,
_dispatch_continuation_invoke_inline(dispatch_object_t dou,
dispatch_invoke_flags_t flags, dispatch_queue_class_t dqu)
{
dispatch_continuation_t dc = dou._dc, dc1;
dispatch_invoke_with_autoreleasepool(flags, {
uintptr_t dc_flags = dc->dc_flags;
// Add the item back to the cache before calling the function. This
// allows the 'hot' continuation to be used for a quick callback.
//
// The ccache version is per-thread.
// Therefore, the object has not been reused yet.
// This generates better assembly.
_dispatch_continuation_voucher_adopt(dc, dc_flags);
if (!(dc_flags & DC_FLAG_NO_INTROSPECTION)) {
_dispatch_trace_item_pop(dqu, dou);
}
if (dc_flags & DC_FLAG_CONSUME) {
dc1 = _dispatch_continuation_free_cacheonly(dc);
} else {
dc1 = NULL;
}
if (unlikely(dc_flags & DC_FLAG_GROUP_ASYNC)) {
_dispatch_continuation_with_group_invoke(dc);
} else {
_dispatch_client_callout(dc->dc_ctxt, dc->dc_func);
_dispatch_trace_item_complete(dc);
}
if (unlikely(dc1)) {
_dispatch_continuation_free_to_cache_limit(dc1);
}
});
_dispatch_perfmon_workitem_inc();
}
static inline void
_dispatch_continuation_with_group_invoke(dispatch_continuation_t dc)
{
struct dispatch_object_s *dou = dc->dc_data;
unsigned long type = dx_type(dou);
if (type == DISPATCH_GROUP_TYPE) {
_dispatch_client_callout(dc->dc_ctxt, dc->dc_func);
_dispatch_trace_item_complete(dc);
dispatch_group_leave((dispatch_group_t)dou);
} else {
DISPATCH_INTERNAL_CRASH(dx_type(dou), "Unexpected object type");
}
}
在这里就会根据 flag 中是否有 DC_FLAG_GROUP_ASYNC 来判断是否需要执行 dispatch_group_leave,其余的依旧是基本一致。下面就来看下在 enter 和 leave 中做了什么。
进入
void
dispatch_group_enter(dispatch_group_t dg)
{
// The value is decremented on a 32bits wide atomic so that the carry
// for the 0 -> -1 transition is not propagated to the upper 32bits.
uint32_t old_bits = os_atomic_sub_orig2o(dg, dg_bits,
DISPATCH_GROUP_VALUE_INTERVAL, acquire);
uint32_t old_value = old_bits & DISPATCH_GROUP_VALUE_MASK;
if (unlikely(old_value == 0)) {
_dispatch_retain(dg); // <rdar://problem/22318411>
}
if (unlikely(old_value == DISPATCH_GROUP_VALUE_MAX)) {
DISPATCH_CLIENT_CRASH(old_bits,
"Too many nested calls to dispatch_group_enter()");
}
}
enter 很简单,就是更新一下 dg_bits 记录当前 enter 的数量,但是 enter 的数量是有一个最大值的,就是 DISPATCH_GROUP_VALUE_MAX,这里简单理解下 dg_bits,刚创建的 group dg_bits 的值一般是 0,然后每一次 enter 就将其减去 DISPATCH_GROUP_VALUE_INTERVAL,所以正常情况下 dg_bits 的值应该一直小于 0,但是这里判断,当 old_value 等于 DISPATCH_GROUP_VALUE_INTERVAL 这样一个正数的时候报错有太多的 enter,应该是当 dg_bits 减太多后,溢出了变成了 DISPATCH_GROUP_VALUE_INTERVAL。
退出
void
dispatch_group_leave(dispatch_group_t dg)
{
// The value is incremented on a 64bits wide atomic so that the carry for
// the -1 -> 0 transition increments the generation atomically.
uint64_t new_state, old_state = os_atomic_add_orig2o(dg, dg_state,
DISPATCH_GROUP_VALUE_INTERVAL, release);
uint32_t old_value = (uint32_t)(old_state & DISPATCH_GROUP_VALUE_MASK);
if (unlikely(old_value == DISPATCH_GROUP_VALUE_1)) {
old_state += DISPATCH_GROUP_VALUE_INTERVAL;
do {
new_state = old_state;
if ((old_state & DISPATCH_GROUP_VALUE_MASK) == 0) {
new_state &= ~DISPATCH_GROUP_HAS_WAITERS;
new_state &= ~DISPATCH_GROUP_HAS_NOTIFS;
} else {
// If the group was entered again since the atomic_add above,
// we can't clear the waiters bit anymore as we don't know for
// which generation the waiters are for
new_state &= ~DISPATCH_GROUP_HAS_NOTIFS;
}
if (old_state == new_state) break;
} while (unlikely(!os_atomic_cmpxchgv2o(dg, dg_state,
old_state, new_state, &old_state, relaxed)));
return _dispatch_group_wake(dg, old_state, true);
}
if (unlikely(old_value == 0)) {
DISPATCH_CLIENT_CRASH((uintptr_t)old_value,
"Unbalanced call to dispatch_group_leave()");
}
}
与进入相反,退出就是给 dg_bits 加上 DISPATCH_GROUP_VALUE_INTERVAL,然后判断此时的值是多少,由于当 enter 和 leave 一致时的 dg_bits 是 0,所以当此时 dg_bits 再次来到 0 时,一定是 enter 和 leave 成对出现了,所以就会重置状态,并调用 _dispatch_group_wake 唤醒正在等待的任务。
static void
_dispatch_group_wake(dispatch_group_t dg, uint64_t dg_state, bool needs_release)
{
uint16_t refs = needs_release ? 1 : 0; // <rdar://problem/22318411>
if (dg_state & DISPATCH_GROUP_HAS_NOTIFS) {
dispatch_continuation_t dc, next_dc, tail;
// Snapshot before anything is notified/woken <rdar://problem/8554546>
dc = os_mpsc_capture_snapshot(os_mpsc(dg, dg_notify), &tail);
do {
dispatch_queue_t dsn_queue = (dispatch_queue_t)dc->dc_data;
next_dc = os_mpsc_pop_snapshot_head(dc, tail, do_next);
_dispatch_continuation_async(dsn_queue, dc,
_dispatch_qos_from_pp(dc->dc_priority), dc->dc_flags);
_dispatch_release(dsn_queue);
} while ((dc = next_dc));
refs++;
}
if (dg_state & DISPATCH_GROUP_HAS_WAITERS) {
_dispatch_wake_by_address(&dg->dg_gen);
}
if (refs) _dispatch_release_n(dg, refs);
}
首先 dispatch group 有两种功能,分别是设置 notify 和设置 waiter,notify 的功能是给 dispatch group 添加一个等待任务,当 group 判断自己已经退出之后就去执行这些任务,waiter 的功能时调用的线程会阻塞,直到 group 执行完。而上面的函数中就可以看到,对于 notify,直接通过 _dispatch_continuation_async 将这个任务添加到对应的队列中就完了,所以从这里可以看出,我们通过 notify 设置的回调,不是 group 执行完就可以立刻执行,此时只是把任务提交到队列中,具体还要看队列自己的调度。而 waiter 线程的唤醒则是通过 _dispatch_wake_by_address 完成,具体的到下面在看。
监听
dispatch group 提供的监听功能就是这一对 wait/notify 函数,俩函数的功能上面也刚说过,下面就两种功能的实现再分析下。
wait
intptr_t
dispatch_group_wait(dispatch_group_t dg, dispatch_time_t timeout)
{
uint64_t old_state, new_state;
os_atomic_rmw_loop2o(dg, dg_state, old_state, new_state, relaxed, {
if ((old_state & DISPATCH_GROUP_VALUE_MASK) == 0) {
os_atomic_rmw_loop_give_up_with_fence(acquire, return 0);
}
if (unlikely(timeout == 0)) {
os_atomic_rmw_loop_give_up(return _DSEMA4_TIMEOUT());
}
new_state = old_state | DISPATCH_GROUP_HAS_WAITERS;
if (unlikely(old_state & DISPATCH_GROUP_HAS_WAITERS)) {
os_atomic_rmw_loop_give_up(break);
}
});
return _dispatch_group_wait_slow(dg, _dg_state_gen(new_state), timeout);
}
static intptr_t
_dispatch_group_wait_slow(dispatch_group_t dg, uint32_t gen,
dispatch_time_t timeout)
{
for (;;) {
int rc = _dispatch_wait_on_address(&dg->dg_gen, gen, timeout, 0);
if (likely(gen != os_atomic_load2o(dg, dg_gen, acquire))) {
return 0;
}
if (rc == ETIMEDOUT) {
return _DSEMA4_TIMEOUT();
}
}
}
首先 fast 模式,先判断 group 是不是处于未 enter 状态,如果是直接返回。如果是 enter 状态的话,就标记下有线程在等待 group 完成,然后走 slow 模式,调用 _dispatch_group_wait_slow 进行线程阻塞。
在这个函数中,它是用 _dispatch_wait_on_address 进行对地址的监听,应该是通过监听这块地址的值有没有发生改变,具体的实现不清楚,但不妨碍。剩下的就是等待了。然后在 _dispatch_group_wake 中进行唤醒的时候,就是调用 _dispatch_wake_by_address 来唤醒监听这个地址的线程,从而完成一个闭环。
notify
static inline void
_dispatch_group_notify(dispatch_group_t dg, dispatch_queue_t dq,
dispatch_continuation_t dsn)
{
uint64_t old_state, new_state;
dispatch_continuation_t prev;
dsn->dc_data = dq;
_dispatch_retain(dq);
prev = os_mpsc_push_update_tail(os_mpsc(dg, dg_notify), dsn, do_next);
if (os_mpsc_push_was_empty(prev)) _dispatch_retain(dg);
os_mpsc_push_update_prev(os_mpsc(dg, dg_notify), prev, dsn, do_next);
if (os_mpsc_push_was_empty(prev)) {
os_atomic_rmw_loop2o(dg, dg_state, old_state, new_state, release, {
new_state = old_state | DISPATCH_GROUP_HAS_NOTIFS;
if ((uint32_t)old_state == 0) {
os_atomic_rmw_loop_give_up({
return _dispatch_group_wake(dg, new_state, false);
});
}
});
}
}
notify 函数的话,就是先将等待执行的任务添加到 dispatch group 的 dg_notify 列表中,然后判断 group 是否是 enter 状态,如果是的话就直接调用 _dispatch_group_wake 去执行了,如果不是就啥也不做,等 group 自己调用 wake 就完了。
从这里可以看出一点不同,wait 的时候是先判断状态,再走 slow 模式。notify 则是直接先加到链表中,再判断是否可以直接执行。按理说它也可以先判断状态是否可以直接执行,不行了再添加到链表中的,这样的效率不应该要比现在的实现效率高吗?这里简单思考下的话,由于 libdispatch 中的这些函数都需要是线程安全的,所以如果按照先判断再添加的逻辑,判断这一步是没问题的,问题在于添加,从目前的代码实现来看,这个代码使用了尽量少的原子操作来实现线程安全,连添加结点这一步都不是完全原子的,而如果是在新的这种逻辑下,如果是在一开始就判断了状态,那就需要后面的添加结点的操作都是原子的,否则如果正好在判断状态和添加结点的空档,另一个线程完成了唤醒阶段,那新加的这个 notify 就无法被执行了,得一直到下一个 wake 出现时才能执行,而如果添加结点整个操作都是原子的,相比之下效率就会受到影响,故而,libdispatch 采取了现在的这种方式。