这是一个创建TransformBlock 的方法,它可以防止具有相同密钥的消息并发执行。每条消息的密钥是通过调用提供的keySelector 函数获得的。具有相同密钥的消息相互顺序处理(而不是并行处理)。密钥也作为参数传递给transform 函数,因为它在某些情况下很有用。
public static TransformBlock<TInput, TOutput>
CreateExclusivePerKeyTransformBlock<TInput, TKey, TOutput>(
Func<TInput, TKey, Task<TOutput>> transform,
ExecutionDataflowBlockOptions dataflowBlockOptions,
Func<TInput, TKey> keySelector,
IEqualityComparer<TKey> keyComparer = null)
{
if (transform == null) throw new ArgumentNullException(nameof(transform));
if (keySelector == null) throw new ArgumentNullException(nameof(keySelector));
if (dataflowBlockOptions == null)
throw new ArgumentNullException(nameof(dataflowBlockOptions));
keyComparer = keyComparer ?? EqualityComparer<TKey>.Default;
var internalCTS = CancellationTokenSource
.CreateLinkedTokenSource(dataflowBlockOptions.CancellationToken);
var maxDOP = dataflowBlockOptions.MaxDegreeOfParallelism;
var taskScheduler = dataflowBlockOptions.TaskScheduler;
var perKeySemaphores = new ConcurrentDictionary<TKey, SemaphoreSlim>(
keyComparer);
SemaphoreSlim maxDopSemaphore;
if (maxDOP == DataflowBlockOptions.Unbounded)
{
maxDopSemaphore = new SemaphoreSlim(Int32.MaxValue);
}
else
{
maxDopSemaphore = new SemaphoreSlim(maxDOP, maxDOP);
// The degree of parallelism is controlled by the semaphore
dataflowBlockOptions.MaxDegreeOfParallelism = DataflowBlockOptions.Unbounded;
// Use a limited-concurrency scheduler for preserving the processing order
dataflowBlockOptions.TaskScheduler = new ConcurrentExclusiveSchedulerPair(
taskScheduler, maxDOP).ConcurrentScheduler;
}
var block = new TransformBlock<TInput, TOutput>(async item =>
{
var key = keySelector(item);
var perKeySemaphore = perKeySemaphores
.GetOrAdd(key, _ => new SemaphoreSlim(1, 1));
// Continue on captured context before invoking the transform
await perKeySemaphore.WaitAsync(internalCTS.Token);
try
{
await maxDopSemaphore.WaitAsync(internalCTS.Token);
try
{
return await transform(item, key).ConfigureAwait(false);
}
catch (Exception ex) when (!(ex is OperationCanceledException))
{
internalCTS.Cancel(); // The block has failed
throw;
}
finally
{
maxDopSemaphore.Release();
}
}
finally
{
perKeySemaphore.Release();
}
}, dataflowBlockOptions);
dataflowBlockOptions.MaxDegreeOfParallelism = maxDOP; // Restore initial value
dataflowBlockOptions.TaskScheduler = taskScheduler; // Restore initial value
return block;
}
使用示例:
var validator = CreateExclusivePerKeyTransformBlock<Uri, string, bool>(
async (uri, host) =>
{
return (await _httpClient.GetAsync(uri, HttpCompletionOption
.ResponseHeadersRead, token)).IsSuccessStatusCode;
},
new ExecutionDataflowBlockOptions
{
MaxDegreeOfParallelism = 30,
CancellationToken = token,
},
keySelector: uri => uri.Host,
keyComparer: StringComparer.OrdinalIgnoreCase);
支持所有execution options(MaxDegreeOfParallelism、BoundedCapacity、CancellationToken、EnsureOrdered 等)。
下面是接受同步委托的CreateExclusivePerKeyTransformBlock 的重载,以及另一个返回ActionBlock 而不是TransformBlock 的方法+重载,具有相同的行为。
public static TransformBlock<TInput, TOutput>
CreateExclusivePerKeyTransformBlock<TInput, TKey, TOutput>(
Func<TInput, TKey, TOutput> transform,
ExecutionDataflowBlockOptions dataflowBlockOptions,
Func<TInput, TKey> keySelector,
IEqualityComparer<TKey> keyComparer = null)
{
if (transform == null) throw new ArgumentNullException(nameof(transform));
return CreateExclusivePerKeyTransformBlock(
(item, key) => Task.FromResult(transform(item, key)),
dataflowBlockOptions, keySelector, keyComparer);
}
// An ITargetBlock is similar to an ActionBlock
public static ITargetBlock<TInput>
CreateExclusivePerKeyActionBlock<TInput, TKey>(
Func<TInput, TKey, Task> action,
ExecutionDataflowBlockOptions dataflowBlockOptions,
Func<TInput, TKey> keySelector,
IEqualityComparer<TKey> keyComparer = null)
{
if (action == null) throw new ArgumentNullException(nameof(action));
var block = CreateExclusivePerKeyTransformBlock(async (item, key) =>
{ await action(item, key).ConfigureAwait(false); return (object)null; },
dataflowBlockOptions, keySelector, keyComparer);
block.LinkTo(DataflowBlock.NullTarget<object>());
return block;
}
public static ITargetBlock<TInput>
CreateExclusivePerKeyActionBlock<TInput, TKey>(
Action<TInput, TKey> action,
ExecutionDataflowBlockOptions dataflowBlockOptions,
Func<TInput, TKey> keySelector,
IEqualityComparer<TKey> keyComparer = null)
{
if (action == null) throw new ArgumentNullException(nameof(action));
return CreateExclusivePerKeyActionBlock(
(item, key) => { action(item, key); return Task.CompletedTask; },
dataflowBlockOptions, keySelector, keyComparer);
}
警告:这个类为每个键分配一个SemaphoreSlim,并保持对它的引用,直到类实例最终被垃圾回收。如果不同键的数量很大,这可能是一个问题。有一个分配较少的异步锁here 的实现,它在内部仅存储当前正在使用的SemaphoreSlims(加上一小部分可以重用的已释放SemaphoreSlims),它可以替换@此实现使用 987654341@。