You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

是否存在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 09:28:33