它解决的不是“再造一个 Channel”#

📅 2026/7/22 3:36:54 👁️ 阅读次数 📝 编程学习
它解决的不是“再造一个 Channel”#

ConcurrentQueue、BlockingCollection 和 Channel 都是非常实用的基础组件。

特别是 Channel,它已经覆盖了高性能异步生产消费、等待唤醒、背压和有界容量等大量通用场景。对于简单的单队列生产消费模型,它通常已经足够好用。

ConcurrentQueue 更适合直接的线程安全 FIFO 操作,BlockingCollection 则在生产消费模型中提供阻塞等待和有界容量。它们都是可靠的基础组件,只是抽象层次更接近单个队列。

BufferQueue 的关注点不同。它把实际业务中经常需要自己补齐的一组能力放到同一个模型里:

使用 (T, TopicName) 隔离不同类型、不同用途的数据;
一个 Topic 可以划分为多个 Partition;
不同 Consumer Group 各自维护独立的消费进度;
同一个 Group 内的多个 Consumer 分摊 Partition;
原生批量拉取,减少逐条消费带来的调度与同步开销;
同时支持 Pull 和 Push 两种消费方式;
同时支持 Auto Commit 和 Manual Commit;
Memory 模式支持有界容量;
MemoryMappedFile 模式支持数据、生产进度和消费进度的本地持久化。
它适合“业务和队列在同一个进程内,但消费模型已经比单条 FIFO 更复杂”的场景,而不是用来替代 Kafka、RabbitMQ 或 RocketMQ 这类跨进程、跨机器的消息系统。

Benchmark:优势主要在批量消费#
项目使用 BenchmarkDotNet,对比 BufferQueue 在 Memory 模式下与 Channel、BlockingCollection 的并发生产和消费表现。下面重点展示与 Channel 的代表性结果。

测试消息均为 int,取值从 0 到 8191,共 8192 条;这里的 MessageSize 表示消息条数,不是单条消息的字节大小。生产测试把这组整数分片给 12 个并发任务,消费测试在计时前把同一组整数预先写入队列。BufferQueue 使用 12 个 Partition,表中数据运行于 Apple M2 Max 和 .NET 10。

生产性能#
模式 MessageSize Producers Channel Mean BufferQueue(Memory)Mean 结果
Unbounded 8192 12 287.0 μs 335.0 μs Channel 约快 1.17x
Bounded 8192 12 300.8 μs 364.1 μs Channel 约快 1.21x
消费性能#
模式 MessageSize BatchSize Consumers Channel Mean BufferQueue(Memory)Mean 结果
Unbounded 8192 1 12 3,146.52 μs 815.80 μs 该组参数下 BufferQueue 约快 3.9x
Bounded 8192 1 12 2,118.13 μs 750.73 μs 该组参数下 BufferQueue 约快 2.8x
Unbounded 8192 100 12 3,384.68 μs 49.25 μs 该组参数下 BufferQueue 约快 69x
Bounded 8192 100 12 2,158.57 μs 53.95 μs 该组参数下 BufferQueue 约快 40x
Unbounded 8192 1000 12 3,485.82 μs 33.97 μs 该组参数下 BufferQueue 约快 103x
Bounded 8192 1000 12 2,115.11 μs 35.68 μs 该组参数下 BufferQueue 约快 59x
这些结果表达了一个很明确的定位:BufferQueue 的生产性能已经接近 Channel,但它真正突出的方向是批量消费吞吐。批次越大,逐条读取、同步和交付成本被摊薄得越明显。

性能数字必须结合工作负载理解。上面的倍数只代表这组消息数量、并发度和 Batch Size,不意味着任意业务都能得到同样结果;当前消费场景中 BufferQueue 的内存分配也高于 Channel。IterationSetup 中的预填充不计入测量,测试衡量的是并发排空队列与批次交付开销,不包含生产和逐条业务处理;Channel 使用 TryRead 逐条读取,BufferQueue 启用 Auto Commit,BatchSize 只影响 BufferQueue。BatchSize 表示每批上限;在 8192 条数据和 12 个 Consumers 下,1000 这一档通常是每个 Consumer 用一批读完分配到的约 682 或 683 条数据。单条处理成本上升后,端到端差距会缩小,因此这些结果更适合用于观察趋势。项目提供了完整的 Benchmark 源码,可以在目标机器上用真实的数据类型和批次大小重新测试:

dotnet run -c Release --project tests/BufferQueue.Benchmarks/BufferQueue.Benchmarks.csproj
先看核心模型#
BufferQueue 中最重要的四个概念是 Topic、Partition、Consumer Group 和 Consumer。

BufferQueue 中按类型和 Topic 隔离的队列

同一个类型可以注册多个 Topic,同一个 Topic 名也可以承载不同类型。每一对 (T, TopicName) 都对应一个独立的强类型队列。例如:

(OrderCreated, “order-events”)
(OrderCancelled, “order-events”)
(OrderCreated, “audit-events”)
在公共 API 和 DI 模型中,这三个组合是不同的队列;Memory 模式下它们也完全独立。需要注意的是,MemoryMappedFile 的磁盘路径不包含类型名:同一个 DataDirectory 下,TopicName 必须跨消息类型保持唯一。如果不同类型需要复用同一个名称,应为它们配置不同的 DataDirectory。

每个 Topic 内部可以包含多个 Partition。Producer 使用轮询方式把数据分发到各个 Partition;创建一个 Consumer Group 时,这些 Partition 会被均分给组内的 Consumer。

Partition 与 Consumer Group 的关系

这里有两个很关键的语义:

不同 Consumer Group 会在当前可读范围内独立消费数据,并分别推进自己的进度。
同一个 Group 内,一条 Partition 只会分配给一个 Consumer,用分区分配实现负载均衡。
例如,一个订单事件 Topic 可以同时被“更新库存”和“生成报表”两个 Group 消费;两个业务互不抢占数据,也不会互相覆盖进度。

两种存储模式怎么选#
维度 Memory MemoryMappedFile
数据位置 当前进程内存 本地内存映射分段文件
重启后恢复 不支持 支持已到达持久化边界的数据和已提交进度
序列化 不需要 目前内置 System.Text.Json、MessagePack 和 unmanaged struct 序列化器,也支持自定义实现
容量控制 支持 BoundedCapacity 当前不支持有界容量
Segment 清理 所有 Group 越过后复用内存 Segment 可按最慢 Group 的已提交进度删除完整文件 Segment
典型场景 可丢弃的实时采样、临时聚合与批处理 希望应用重启后继续消费的本地缓冲
如果应用重启时丢失少量待处理数据可以接受,或者数据能够重新生成,并且重点是延迟和消费吞吐,选择 Memory。

如果数据和已提交进度需要跨重启保留,同时仍然只运行一个 active queue 实例,选择 MemoryMappedFile。

快速接入#
核心包和可选的 MemoryMappedFile 扩展都已经发布到 NuGet:

dotnet add package BufferQueue
dotnet add package BufferQueue.MemoryMappedFile
其中 BufferQueue 包含公共队列模型、Memory 存储和 Push Consumer;BufferQueue.MemoryMappedFile 是可选的持久化扩展,并依赖核心包。

下面在同一个应用中注册两个 Topic:允许少量丢失、可以重新采集的非关键遥测数据使用 Memory,需要跨重启保留的订单事件使用 MemoryMappedFile。

using BufferQueue;
using BufferQueue.Memory;
using BufferQueue.MemoryMappedFile;
using BufferQueue.PushConsumer;

builder.Services.AddBufferQueue(bufferQueue =>
{
bufferQueue
.UseMemory(memory =>
{
memory.AddTopic(options =>
{
options.TopicName = “telemetry-samples”;
options.PartitionNumber = 4;
options.BoundedCapacity = 100_000;
});
})
.UseMemoryMappedFile(memoryMappedFile =>
{
memoryMappedFile.AddTopic(options =>
{
options.TopicName = “order-events”;
options.PartitionNumber = 4;
options.DataDirectory = “/var/lib/bufferqueue”;
});
})
.AddPushCustomers(typeof(Program).Assembly);
});
每个 (T, TopicName) 只应注册到一种存储模式。Memory 和 MemoryMappedFile 可以共存,但不要把同一个组合重复注册到两种模式;使用 MemoryMappedFile 时,同一 DataDirectory 内还要确保不同类型不会复用相同的 TopicName。

生产数据#
Topic 在依赖声明时已经确定,可以直接使用 .NET 的 Keyed Service 注入 Producer:

using Microsoft.Extensions.DependencyInjection;

public sealed class OrderService(
[FromKeyedServices(“order-events”)]
IBufferProducer producer)
{
public ValueTask PublishAsync(OrderCreated order) =>
producer.ProduceAsync(order);
}
如果 Topic 需要在运行时决定,也可以通过 IBufferQueue 获取:

var producer = bufferQueue.GetProducer(“order-events”);
await producer.ProduceAsync(orderCreated);
Memory Topic 设置 BoundedCapacity 后,队列满时 ProduceAsync 会抛出 MemoryBufferQueueFullException。不希望使用异常控制流程时,可以改用 TryProduceAsync:

var telemetryProducer = bufferQueue.GetProducer(“telemetry-samples”);
var accepted = await telemetryProducer.TryProduceAsync(sample);
if (!accepted)
{
// 当前遥测采样可以按业务约定丢弃或降级处理
}
Pull 模式批量消费#
下面创建 4 个 Consumer,并把 order-events 的 4 个 Partition 分配给它们。这里关闭 Auto Commit,业务处理成功后再手动提交:

var consumers = bufferQueue.CreatePullConsumers(
new BufferPullConsumerOptions
{
TopicName = “order-events”,
GroupName = “inventory-update”,
BatchSize = 200,
AutoCommit = false
},
consumerNumber: 4);

await Task.WhenAll(
consumers.Select(consumer => ConsumeAsync(consumer, stoppingToken)));

static async Task ConsumeAsync(
IBufferPullConsumer consumer,
CancellationToken cancellationToken)
{
await foreach (var batch in consumer.ConsumeAsync(cancellationToken))
{
await UpdateInventoryAsync(batch, cancellationToken);
await consumer.CommitAsync();
}
}
单个 Consumer 会顺序推进自己负责的 Partition。要提高组内并行度,应在创建 Group 时增加 Consumer 数量;Consumer 数量不能超过 Partition 数量,并且当前不支持在运行时动态增减同一 Group 的 Consumer。

GroupName 建议按消费用途命名,例如 inventory-update、sales-reporting。执行同一用途的多个 Consumer 共享一个 Group,不同业务用途使用不同 Group;不要把机器名、进程号或实例编号作为 Group 名的一部分。

Push 模式消费#
如果不想自己管理消费循环,可以扫描带有 BufferPushCustomerAttribute 的类型,由 Hosted Service 驱动消费:

上面的注册代码已经通过 .AddPushCustomers(typeof(Program).Assembly) 开启扫描。BufferPushCustomerAttribute 是当前 API 的名称,正文仍统一称它创建的消费者为 Push Consumer。

using BufferQueue.PushConsumer;
using Microsoft.Extensions.DependencyInjection;

[BufferPushCustomer(
topicName: “telemetry-samples”,
groupName: “telemetry-storage”,
batchSize: 500,
serviceLifetime: ServiceLifetime.Singleton,
concurrency: 4)]
public sealed class TelemetrySampleConsumer : IBufferAutoCommitPushConsumer
{
public Task ConsumeAsync(
IEnumerable batch,
CancellationToken cancellationToken) =>
WriteTelemetryAsync(batch, cancellationToken);
}
Push Consumer 仍然建立在同一套 Pull Consumer 模型之上,所以 Partition 分配、批量大小和提交语义保持一致。

Auto Commit 和 Manual Commit 到底有什么区别#
这个选择直接决定失败时的行为。

Auto Commit#
一次 Pull 成功后,BufferQueue 会立即推进消费进度,然后把批次交给业务代码。它使用方便、吞吐路径短,但业务处理失败时,这个批次不会因为“没有提交”而自动重放。

适合允许少量数据遗漏、数据可以重新生成,或者业务自身有独立恢复机制的场景。

Manual Commit#
业务处理成功后显式调用 CommitAsync。如果处理失败且没有提交,同一批数据可以再次被读取,提供 at-least-once 语义。

这也意味着业务处理应尽量做到幂等,因为“至少一次”允许重复,而不承诺“恰好一次”。

未提交只保证消费进度不会推进,BufferQueue 不会主动调度失败重试。应用需要继续或重新启动消费循环,才能再次读取这批数据。

在 MemoryMappedFile 模式下,Commit 还会先强制待提交日志到达 flush 边界,再持久化 Consumer Offset,使已到达持久化边界的记录和已提交进度可以跨进程重启恢复。

Memory 模式为什么适合批量消费#
Memory Partition 不是一个不断扩容的大数组,而是由固定大小的 Segment 组成的链表:

head segment -> segment -> … -> tail segment
Producer 向尾部 Segment 追加数据。写满后,Partition 会创建新 Segment,或者复用已经被所有 Consumer Group 完整消费的旧 Segment。

这里的“所有 Group”很重要。假设报表 Group 比实时计算 Group 慢很多,只要报表 Group 还没有越过某个 Segment,这段内存就不能被复用。这样既能支持多个独立进度,又不会覆盖慢消费者尚未读取的数据。

并发方面,Producer 支持多线程调用。Memory Queue 会在一段很短的临界区内完成 Partition 轮询、容量计数和追加;Consumer 读取已经发布的数据时不获取这个 append lock。多个 Group 可以独立读取和推进各自的进度。

批量消费则把多条数据一次性返回给业务代码,减少逐条读取带来的状态切换、同步和调度成本。这也是 BufferQueue 相比通用 Queue 最明显的性能侧重点。

MemoryMappedFile:把同一套消费模型带到磁盘#
MemoryMappedFile 模式没有在上层重新实现一套队列。两种存储共享 Topic、Producer、Consumer、Consumer Group、Partition 分配、等待唤醒和提交逻辑,差异被隔离在内部 Partition 抽象之后:

Application
|
v
IBufferQueue / IBufferProducer
|
v
BufferQueue
|
±- BufferPullConsumer
±- IBufferPartition[]
|
±- MemoryBufferPartition
±- MemoryMappedFileBufferPartition
配置示例#
MemoryMappedFile 的所有参数都按 Topic 配置。下面两种写法分别展示 MessagePack 和 unmanaged struct 的配置方式;它们是互斥示例,同一个 (T, TopicName) 不要重复注册。

MessagePack、批量 Flush 与 Segment 清理策略#
如果消息量较大,希望减少 Serializer 和显式 Flush 的开销,可以给持久化类型定义稳定的 MessagePack Contract:

dotnet add package MessagePack
using MessagePack;

[MessagePackObject]
public sealed class OrderCreated
{
[Key(0)]
public long OrderId { get; set; }

[Key(1)] public long CustomerId { get; set; } [Key(2)] public decimal Amount { get; set; } [Key(3)] public long CreatedAtUnixTimeMilliseconds { get; set; }

}
然后配置 MessagePack Serializer、批量 Flush 与 Segment 清理策略:

using System.IO;
using BufferQueue;
using BufferQueue.MemoryMappedFile;

builder.Services.AddBufferQueue(bufferQueue =>
{
bufferQueue.UseMemoryMappedFile(memoryMappedFile =>
{
memoryMappedFile.AddTopic(options =>
{
options.TopicName = “order-events”;
options.PartitionNumber = 4;
options.DataDirectory = Path.Combine(
builder.Environment.ContentRootPath,
“bufferqueue-data”);
options.SegmentSizeInBytes = 64L * 1024 * 1024;
options.MaxRetainedConsumedSegments = 2;
options.FlushStrategy = MemoryMappedFileFlushStrategy.Batch;
options.FlushBatchSize = 100;
options.Serializer =
new MessagePackMemoryMappedFileSerializer();
});
});
});