是否存在BufferBlock的FILO实现?双BufferBlock消息路由及消费优先级问询
嘿,这个场景用TPL Dataflow完全能搞定,我来给你捋捋具体的实现思路和代码示例,应该能帮你快速搭起来:
实现方案
1. 核心组件定义
首先咱们要搞定三个关键块:
- 主缓冲区P:用普通的
BufferBlock<T>,设置好BoundedCapacity达到你要的消息上限,它默认是FIFO顺序,正好符合需求。 - 备用FILO缓冲区S:TPL Dataflow没有原生的FILO BufferBlock,不过咱们可以自己封装一个基于栈的块,保证消息先进后出。
- 路由块:负责把输入消息先往P塞,P满了就转去S。
2. 自定义FILO缓冲区块
先写个简单的FiloBufferBlock<T>,内部用ConcurrentStack存消息,消费的时候弹栈顶元素,完美实现FILO:
using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Threading.Tasks; using System.Threading.Tasks.Dataflow; public class FiloBufferBlock<T> : ITargetBlock<T>, ISourceBlock<T> { private readonly ConcurrentStack<T> _messageStack = new ConcurrentStack<T>(); private readonly TaskCompletionSource<bool> _completionSource = new TaskCompletionSource<bool>(); private readonly object _lockObj = new object(); private List<ITargetBlock<T>> _linkedTargets = new List<ITargetBlock<T>>(); public Task Completion => _completionSource.Task; public void Complete() { lock (_lockObj) { _completionSource.TrySetResult(true); } } public void Fault(Exception error) { lock (_lockObj) { _completionSource.TrySetException(error); } } public IDisposable LinkTo(ITargetBlock<T> target, DataflowLinkOptions linkOptions) { lock (_lockObj) { _linkedTargets.Add(target); TrySendPendingMessages(); } return new Unlinker(this, target); } private bool TrySendPendingMessages() { while (_messageStack.TryPop(out var message)) { foreach (var target in _linkedTargets.ToList()) { if (target.Post(message)) { return true; } } _messageStack.Push(message); return false; } return false; } DataflowMessageStatus ITargetBlock<T>.OfferMessage(DataflowMessageHeader header, T value, ISourceBlock<T> source, bool consumeToAccept) { if (consumeToAccept && source != null) { source.ConsumeMessage(header, this, out _); } lock (_lockObj) { if (_completionSource.Task.IsCompleted) { return DataflowMessageStatus.DecliningPermanently; } _messageStack.Push(value); TrySendPendingMessages(); return DataflowMessageStatus.Accepted; } } T ISourceBlock<T>.ConsumeMessage(DataflowMessageHeader header, ITargetBlock<T> target, out bool consumed) { consumed = false; throw new NotImplementedException(); } bool ISourceBlock<T>.ReserveMessage(DataflowMessageHeader header, ITargetBlock<T> target) { throw new NotImplementedException(); } void ISourceBlock<T>.ReleaseReservation(DataflowMessageHeader header, ITargetBlock<T> target) { throw new NotImplementedException(); } private class Unlinker : IDisposable { private readonly FiloBufferBlock<T> _block; private readonly ITargetBlock<T> _target; public Unlinker(FiloBufferBlock<T> block, ITargetBlock<T> target) { _block = block; _target = target; } public void Dispose() { lock (_block._lockObj) { _block._linkedTargets.Remove(_target); } } } }
3. 路由与消费逻辑
接下来搭建整体流程:输入消息先经过路由块,优先往P塞,P满了就去S;消费时优先从S取消息,S空了再处理P的消息。
using System; using System.Threading.Tasks; using System.Threading.Tasks.Dataflow; class Program { static async Task Main(string[] args) { // 配置主缓冲区容量,比如设为5 const int MainBufferCapacity = 5; // 1. 创建主缓冲区P(FIFO,有界) var mainBuffer = new BufferBlock<string>(new DataflowBlockOptions { BoundedCapacity = MainBufferCapacity, EnsureOrdered = true }); // 2. 创建备用FILO缓冲区S var filoBuffer = new FiloBufferBlock<string>(); // 3. 创建路由块:分发消息到P或S var routerBlock = new ActionBlock<string>(msg => { // 尝试发送到主缓冲区,成功则结束,否则转去备用区 bool sentToMain = mainBuffer.Post(msg); if (!sentToMain) { filoBuffer.Post(msg); Console.WriteLine($"Message '{msg}' routed to S (P is full)"); } else { Console.WriteLine($"Message '{msg}' routed to P"); } }); // 4. 实现优先消费逻辑:先S后P async Task ConsumeMessages() { while (true) { // 先尝试从S取消息(FILO) if (filoBuffer.TryReceive(out var sMsg)) { Console.WriteLine($"Consumed from S: {sMsg}"); } else { // S空了,从P取消息 var pMsg = await mainBuffer.ReceiveAsync(); Console.WriteLine($"Consumed from P: {pMsg}"); } // 退出条件:所有块完成且无剩余消息 if (mainBuffer.Completion.IsCompleted && filoBuffer.Completion.IsCompleted && !mainBuffer.TryReceive(out _) && !filoBuffer.TryReceive(out _)) { break; } } } // 测试:发送10条消息 for (int i = 1; i <= 10; i++) { var msg = $"Message {i}"; routerBlock.Post(msg); Console.WriteLine($"Sent: {msg}"); } // 标记路由块完成,触发后续缓冲区的完成逻辑 routerBlock.Complete(); await routerBlock.Completion; mainBuffer.Complete(); filoBuffer.Complete(); // 启动消费任务 await ConsumeMessages(); Console.WriteLine("\nAll messages processed!"); } }
关键细节说明
- 路由判断:用
BufferBlock.Post()的返回值判断主缓冲区是否已满,返回false就说明P已经达到容量上限,直接转发到S。 - FILO实现:自定义的
FiloBufferBlock用栈存储消息,消费时总是取最新进入的消息,完美符合需求。 - 优先消费:消费循环先检查S是否有消息,只有当S空了才会去P取,保证了S的优先级。
内容的提问来源于stack exchange,提问作者BillHaggerty
相关产品推荐
相关产品推荐

