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

C# ConcurrentBag:每新增N个对象时如何安全清空?异步日志批量发送方案

针对你在WCF服务里实现异步CloudWatch日志器的需求,结合多线程场景和AWS SDK的限制,我整理了一套最优的实现方案,兼顾性能、线程安全和日志可靠性:

核心设计思路

首先纠正一个误区:你不需要在发送日志时锁定其他线程的日志添加操作。AWS SDK只是不允许同时发送多个请求,但日志的收集和发送可以完全异步分离——业务线程只管往线程安全的队列里塞日志,后台单独用一个(或受控数量的)线程处理批量发送,这样既不会阻塞业务逻辑,又能满足SDK的并发限制。

另外,推荐用ConcurrentQueue替代ConcurrentBag:ConcurrentBag是为线程本地存储优化的集合,元素顺序不确定,而日志通常需要严格的时序性,ConcurrentQueue的FIFO特性更适合日志场景。

具体实现方案
  1. 线程安全单例日志器:用Lazy<T>实现懒加载的线程安全单例,确保整个WCF服务实例只有一个日志器实例。
  2. 双触发批量发送机制:
    • 当队列中的日志数量达到设定的批量阈值(比如100条)时,自动触发发送
    • 定时触发(比如每10秒),避免低业务量时日志长期积压
  3. 异步发送与错误重试:用AWS SDK的异步API发送日志,处理常见错误(比如序列令牌过期),并将失败的日志重新放回队列,避免丢失。
  4. 非阻塞日志收集:业务线程调用日志方法时,仅执行入队操作,几乎无性能损耗。
代码示例
using Amazon.CloudWatchLogs;
using Amazon.CloudWatchLogs.Model;
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;

public sealed class CloudWatchAsyncLogger
{
    // 线程安全懒加载单例
    private static readonly Lazy<CloudWatchAsyncLogger> _instance = 
        new Lazy<CloudWatchAsyncLogger>(() => new CloudWatchAsyncLogger());
    public static CloudWatchAsyncLogger Instance => _instance.Value;

    private readonly ConcurrentQueue<InputLogEvent> _logQueue = new ConcurrentQueue<InputLogEvent>();
    private readonly AmazonCloudWatchLogsClient _cwClient;
    private readonly string _logGroupName;
    private readonly string _logStreamName;
    private readonly int _batchSize;
    private readonly TimeSpan _flushInterval;
    private readonly SemaphoreSlim _sendSemaphore = new SemaphoreSlim(1, 1);
    private string _nextSequenceToken;
    private Task _backgroundFlushTask;

    private CloudWatchAsyncLogger()
    {
        // 初始化配置,可根据实际需求调整
        _cwClient = new AmazonCloudWatchLogsClient();
        _logGroupName = "WcfServiceLogs";
        _logStreamName = $"WcfLogStream_{Environment.MachineName}";
        _batchSize = 100; // 注意CloudWatch限制:单次最多1000条或1MB
        _flushInterval = TimeSpan.FromSeconds(10);

        // 启动后台定时刷新任务
        _backgroundFlushTask = Task.Run(BackgroundFlushLoop);
    }

    /// <summary>
    /// 业务线程调用的日志方法
    /// </summary>
    public void Log(string message)
    {
        var logEvent = new InputLogEvent
        {
            Timestamp = DateTime.UtcNow,
            Message = message
        };
        _logQueue.Enqueue(logEvent);

        // 达到批量阈值时触发一次发送
        if (_logQueue.Count >= _batchSize)
        {
            TriggerImmediateFlush();
        }
    }

    /// <summary>
    /// 触发立即刷新(非阻塞)
    /// </summary>
    private void TriggerImmediateFlush()
    {
        // 用信号量确保同时只有一个发送任务在执行
        if (_sendSemaphore.Wait(0))
        {
            try
            {
                if (_backgroundFlushTask.IsCompleted)
                {
                    _backgroundFlushTask = Task.Run(BackgroundFlush);
                }
            }
            finally
            {
                _sendSemaphore.Release();
            }
        }
    }

    /// <summary>
    /// 后台定时刷新循环
    /// </summary>
    private async Task BackgroundFlushLoop()
    {
        while (true)
        {
            await Task.Delay(_flushInterval);
            await BackgroundFlush();
        }
    }

    /// <summary>
    /// 执行批量发送逻辑
    /// </summary>
    private async Task BackgroundFlush()
    {
        var batchLogs = new List<InputLogEvent>();

        // 从队列中取出最多_batchSize条日志
        while (batchLogs.Count < _batchSize && _logQueue.TryDequeue(out var logEvent))
        {
            batchLogs.Add(logEvent);
        }

        if (batchLogs.Count == 0) return;

        try
        {
            var request = new PutLogEventsRequest
            {
                LogGroupName = _logGroupName,
                LogStreamName = _logStreamName,
                LogEvents = batchLogs,
                SequenceToken = _nextSequenceToken
            };

            var response = await _cwClient.PutLogEventsAsync(request);
            _nextSequenceToken = response.NextSequenceToken;
        }
        catch (AmazonCloudWatchLogsException ex)
        {
            // 处理序列令牌过期的常见错误
            if (ex.ErrorCode == "InvalidSequenceTokenException")
            {
                // 从错误消息中提取正确的令牌
                _nextSequenceToken = ex.Message.Split(':')[1].Trim();
                // 将失败的日志放回队列,稍后重试
                foreach (var log in batchLogs)
                {
                    _logQueue.Enqueue(log);
                }
                // 立即重试一次
                await BackgroundFlush();
            }
            else
            {
                // 其他CloudWatch错误,比如网络问题,将日志放回队列
                foreach (var log in batchLogs)
                {
                    _logQueue.Enqueue(log);
                }
                // 可选:记录本地错误日志,便于排查
                // File.AppendAllText("CloudWatchLogError.log", $"{DateTime.UtcNow}: {ex.Message}\n");
            }
        }
        catch (Exception ex)
        {
            // 通用错误处理,确保日志不丢失
            foreach (var log in batchLogs)
            {
                _logQueue.Enqueue(log);
            }
            // File.AppendAllText("CloudWatchLogError.log", $"{DateTime.UtcNow}: Unexpected error - {ex.Message}\n");
        }
        finally
        {
            _sendSemaphore.Release();
        }
    }

    /// <summary>
    /// WCF服务关闭时调用,刷新剩余日志
    /// </summary>
    public async Task FlushOnShutdown()
    {
        await BackgroundFlush();
        _cwClient.Dispose();
    }
}
关键注意事项
  1. WCF生命周期集成:在WCF服务的ServiceHost.Opening事件中初始化日志器,Closing事件中调用FlushOnShutdown(),确保服务关闭时所有剩余日志都能发送。
  2. CloudWatch限制:单次PutLogEvents请求最多包含1000条日志,总大小不超过1MB,所以批量大小建议设置在100-500之间,避免触发限制。
  3. 序列令牌管理:CloudWatch要求每次发送都携带上一次的SequenceToken,必须正确处理令牌过期的情况,否则会发送失败。
  4. 日志持久化 fallback:如果遇到长时间网络故障,建议增加本地文件持久化的逻辑,将无法发送的日志写入本地文件,待恢复后再批量上传,避免日志丢失。
  5. 资源释放:确保AWS客户端在服务关闭时正确释放,避免资源泄漏。

内容的提问来源于stack exchange,提问作者Vulegend

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 09:04:40