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

创建固定大小或时间间隔触发的异步事件序列实现方案问询

这个需求其实很常见,本质上就是实现**批量攒批(batching)**的逻辑——要么攒够指定数量,要么到了时间就一次性推送。结合你的单例场景,我给你整理了一套兼顾线程安全和异步效率的实现方案:

核心思路

你需要一个线程安全的缓存容器暂存事件,同时维护两个触发推送的条件:

  1. 缓存的事件数量达到设定阈值(比如10000条)
  2. 距离上次推送的时间间隔达到设定阈值(比如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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 03:42:29