上次说到 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 采取了现在的这种方式。