创建固定大小或时间间隔触发的异步事件序列实现方案问询
这个需求其实很常见,本质上就是实现**批量攒批(batching)**的逻辑——要么攒够指定数量,要么到了时间就一次性推送。结合你的单例场景,我给你整理了一套兼顾线程安全和异步效率的实现方案:
核心思路
你需要一个线程安全的缓存容器暂存事件,同时维护两个触发推送的条件:
- 缓存的事件数量达到设定阈值(比如10000条)
- 距离上次推送的时间间隔达到设定阈值(比如1分钟)
另外还要处理应用 shutdown 时的剩余事件,避免数据丢失。
具体实现步骤
1. 定义攒批配置和核心容器
在你的单例类里,先声明线程安全的队列、并发控制信号量,以及可调整的攒批参数:
public class EventSenderSingleton { // 线程安全队列,用于暂存待发送事件 private readonly ConcurrentQueue<Event> _eventQueue = new(); // 信号量控制并发,避免同时触发多次批量推送 private readonly SemaphoreSlim _sendSemaphore = new(1, 1); // 攒批配置:可根据业务需求灵活调整 private const int BatchSize = 10000; private readonly TimeSpan BatchInterval = TimeSpan.FromMinutes(1); // 定时触发的定时器(异步友好型) private PeriodicTimer? _batchTimer; // 单例构造函数(饿汉式,也可按需改成懒加载) private EventSenderSingleton() { // 启动定时攒批任务 StartBatchTimer(); // 注册进程退出时的清理逻辑,确保剩余事件被发送 AppDomain.CurrentDomain.ProcessExit += OnProcessExit; } // 单例实例 public static EventSenderSingleton Instance { get; } = new(); }
2. 实现定时触发逻辑
用PeriodicTimer替代传统的Timer,它更适配异步场景,不会阻塞主线程:
private void StartBatchTimer() { _batchTimer = new PeriodicTimer(BatchInterval); // 后台启动定时任务,不阻塞主线程 _ = Task.Run(async () => { try { while (await _batchTimer.WaitForNextTickAsync()) { // 到点触发批量发送 await SendBatchIfNeededAsync(); } } catch (OperationCanceledException) { // 定时器被取消,属于正常退出场景 } }); }
3. 修改原SendEventsAsync方法,改为攒批模式
原来的逐个通知逻辑改成将事件加入队列,同时检查是否达到批量阈值:
public async Task SendEventsAsync(IEnumerable<Event> events) { if (events == null || !events.Any()) return; // 将所有事件加入线程安全队列 foreach (var ev in events) { _eventQueue.Enqueue(ev); } // 如果队列数量达到阈值,立即触发批量发送 if (_eventQueue.Count >= BatchSize) { await SendBatchIfNeededAsync(); } }
4. 实现批量发送的核心逻辑
这里要注意线程安全,用信号量避免并发发送,同时一次性取出队列内的所有事件:
private async Task SendBatchIfNeededAsync() { // 尝试获取信号量,如果已有发送任务在执行则直接返回 if (!await _sendSemaphore.WaitAsync(TimeSpan.Zero)) { return; } try { // 一次性取出队列中所有事件 var batchEvents = new List<Event>(); while (_eventQueue.TryDequeue(out var ev)) { batchEvents.Add(ev); } if (batchEvents.Count == 0) return; // 替换成你的订阅者通知逻辑 Console.WriteLine($"批量推送 {batchEvents.Count} 条事件"); // await NotifySubscribersAsync(batchEvents); } catch (Exception ex) { // 这里可添加异常处理,比如日志记录、重试逻辑(如用Polly库) Console.WriteLine($"批量发送失败:{ex.Message}"); } finally { // 无论成功失败,都释放信号量 _sendSemaphore.Release(); } }
5. 处理应用关闭时的剩余事件
确保进程退出前,把队列里没发完的事件全部推送出去:
private async void OnProcessExit(object? sender, EventArgs e) { // 停止定时器 _batchTimer?.Dispose(); // 发送剩余事件 await SendBatchIfNeededAsync(); }
关键细节说明
- 线程安全:
ConcurrentQueue保证多线程下的队列操作安全,SemaphoreSlim避免同时触发多个批量发送任务,防止重复推送。 - 异步友好:全程用
async/await处理异步逻辑,不会阻塞主线程影响应用性能。 - 灵活性:
BatchSize和BatchInterval可以改成从配置文件读取,方便后续动态调整。 - 异常隔离:批量发送的异常单独捕获,不会影响后续的攒批逻辑,必要时可添加重试机制保障可靠性。
内容的提问来源于stack exchange,提问作者Alpha75
相关产品推荐
相关产品推荐

