首页 / 资讯中心 / 文章详情

Serilog消息队列Sink实战:从选型、接入到重复消费

Serilog消息队列Sink实战:从选型、接入到重复消费 ★ FEATURED ARTICLE
Serilog 的消息队列 Sink 我一直想单独挑出来写一篇。前面几篇把控制台、文件、滚动日志这些常规玩法聊得差不多了但实际项目里一旦走到微服务规模和突发流量之后你会发现日志问题根本不是“写到哪”的问题而是“怎么扛住峰值”的问题。这篇以 .net8 为基底把消息队列 Sinks 从选型、配置、消费端落库到重复消费问题完整摊开讲一遍适合已经用上 Serilog、但正在被日志写入瓶颈卡住的团队。去年我手上有过一个 12 节点的订单服务集群高峰期每秒会产生几万条结构化日志那时候我才真正理解什么叫“日志风暴”。当时第一个方案是文件 Sink 直接落盘第二个方案是写数据库。结果一个是磁盘 IO 被打满一个是数据库连接池被拖垮。折腾了一圈之后我最终把消息队列放在应用和日志存储之间这个问题才算彻底解决。这篇文章就围绕这个场景展开讲清楚为什么需要消息队列、Sink 怎么选、怎么接入以及大家最关心的重复消费问题。1. 为什么日志不能直接写存储一次“日志风暴”后的架构调整1.1 文件落盘和数据库直写被压垮的过程先说说那次事故的具体表现。我们的核心接口在高峰期大约有 3 万 QPS每个请求会记录请求开始、参数摘要、业务处理结果、耗时统计、异常堆栈这几类日志单请求平均产生 4 到 5 条结构化日志算下来每秒要写入十几万条。当时用的是 Serilog 文件 Sink按天滚动单个文件限制 100MB。压测到 80% 峰值流量的时候最先出问题的是磁盘 IO 等待时间从正常的不到 10 毫秒直接飙到接近 300 毫秒。日志线程开始阻塞接着业务接口的 P99 延迟从 120 毫秒上涨到 1800 毫秒整个服务就像被自己的日志勒住了脖子。另一个服务用的是数据库直写方案情况更惨。每条日志一次 INSERT看起来单条只要 2 到 3 毫秒但连接池只有 100 个连接日志写入量一大连接就被日志线程占满真正的业务 SQL 反而拿不到连接。数据库侧的锁等待和事务日志膨胀又把写入延迟进一步放大最后连基础监控表的写入都开始超时。这个教训让我彻底明白日志写入不是“把日志存下来”那么简单它本质上是和业务流量强耦合的一类高吞吐写操作必须和业务解耦。我后来把日志链路改成了“应用写队列、消费端慢慢落库”的模型。应用线程把日志事件交给 Serilog 的批量 SinkSink 攒到一定数量后一次性投递到 RabbitMQ数据库写入完全由独立的消费服务来执行。改造后日志写入对业务线程的阻塞几乎可以忽略数据库只需要按自己的节奏消费即可。这个模型带来的第一个收益就是削峰填谷流量高峰时消息队列像蓄水池一样把日志暂时存住消费端保持稳定的写入速率流量回落后队列积压再慢慢消化。整个过程业务无感知磁盘和数据库也不会再被瞬时打爆。1.2 消息队列削峰填谷日志链路里多一层缓冲的真正价值削峰填谷只是表面收益更深层的价值是解耦了“日志产生”和“日志消费”两个环节的速率强约束。之前直写数据库时日志生产者必须适配数据库的写入能力引入消息队列后生产端只需要保证投递成功即可消费端可以根据磁盘、索引、网络情况自由调整批量大小和写入频率。这也意味着日志存储可以随时升级从文件换成 Elasticsearch、ClickHouse 或别的存储引擎只要消费端适配生产端完全不用动。还有一层收益是集中汇聚。微服务架构下几十个服务的日志散落在各自宿主机上排查问题要到每一台机器上翻文件悲观点说运维成本极高。日志先进消息队列所有服务的日志事件汇聚到同一批队列消费端统一处理后进入同一个索引或数据表。这样不仅日志格式可以统一连跨服务的链路追踪也能在查询层直接打通。我们的排查效率因此提升得很明显以前查一个问题要 SSH 到四五台机器现在直接到日志平台按 traceId 搜一次就出来了。1.3 不是所有日志都该进队列适用场景与边界不过我也要泼点冷水消息队列不是日志方案的万能解药。如果你的系统是单体应用、日日志量不大直接文件 Sink 加一个 Seq 服务就够了引入 Kafka 或 RabbitMQ 反而是给自己找运维麻烦。消息队列本身也是基础设施会宕机、会积压、会重复投递没有足够的运维能力撑着出了问题比写文件更难受。适合上消息队列的场景有几个特征第一日志写入峰值远超存储后端稳定写入能力第二服务实例数多日志需要集中检索第三日志还要做二次加工比如打点统计、告警规则匹配、数据清洗。反过来如果只是本地调试、单机运行、对日志实时性要求不高的场景就不要往队列里塞了。另外要特别提醒业务审计类日志、计费日志这类“绝对不能丢”的日志虽然也可以走消息队列但队列端必须开启持久化、生产端要开启发布确认消费端还要手工 ACK任何一环偷懒都可能造成审计数据缺失这在合规上是事故级别的问题。2. Serilog 消息队列 Sink 全家桶选型前必须知道的事2.1 先分清三个问题连接谁、投递到哪、消费端谁来接做选型前我建议先把三个问题在脑子里过一遍Sink 负责连接什么类型的消息系统日志事件投递到哪个主题或路由消费端用什么方式把消息读出来。这三个问题对应着整条日志链路的入口、通道和出口。Serilog 的 Sink 抽象本身很简单核心就是ILogEventSink接口实现Emit(LogEvent logEvent)方法即可。消息队列类 Sink 通常会在内部用PeriodicBatchingSink包装一次把单条日志缓存在内存里达到批量条数或时间阈值后再一次性投递。这个机制是所有批量型 Sink 的共同基础理解了它后面的参数调优就顺理成章。2.2 RabbitMQ生态最成熟也是坑最多的一个Serilog 生态里对 RabbitMQ 的支持算是最成熟的社区维护的Serilog.Sinks.RabbitMQ包已经很长时间了。RabbitMQ 本身的模型对日志场景很友好Exchange 负责路由Queue 负责缓存消费者按自己的节奏取消息。这就意味着日志的“路由”、“缓存”、“消费”三段可以完全独立配置非常灵活。但坑也集中在这里。这个 Sink 的开发活跃度陆陆续续有波动不同版本的配置项名称有细微差异比如ExchangeType、RouteKey、DeliveryMode这些字段在不同版本里大小写、枚举序列化方式不完全一致。我踩过的具体问题是某个旧版本对 RabbitMQ.Client 6.x 的兼容性不好日志投递在长连接场景下会出现偶发断流升级到新版之后才解决。所以用这个包第一件事就是锁版本别随手 “latest” 拉一个然后在一个稳定的版本上做压测验证。RabbitMQ Sink 的另一个特点是默认只负责把消息送到 Exchange不创建 Queue。很多新手以为配了 Sink 消息就进队列了其实 Exchange 如果没有队列绑定消息会被直接丢弃。正确做法是提前在 RabbitMQ 管理端或者消费端启动代码里把 Exchange、Queue 和 Binding 都声明好确保投递时有消费者在收。我习惯在消费端统一声明队列这样生产端只管投递队列的生命周期由真正负责读的人管理职责更清晰。2.3 Kafka 与 Redis Stream吞吐与轻量的两条路线如果你的日志规模到了每天亿级以上Kafka 会是更合适的通道。Serilog 有第三方 Sink 基于 Confluent.Kafka 封装配置上通过bootstrapServers指定集群地址指定topic后按分区投递。Kafka 的优势在于吞吐量极高分区机制天然支持多消费者并行消费而且分区内消息有序对于按时间线排查问题很有帮助。代价是运维复杂度明显上升消费组、offset、broker 扩容每个环节都要懂小团队贸然上 Kafka 容易从“日志问题”变成“Kafka 问题”。如果规模没那么大只是想把十几个实例的日志收拢一下我推荐用 Redis Stream。Redis 在绝大多数团队里已经存在Stream 数据结构从 5.0 开始就稳定了消费者组、待处理消息列表这些机制足够支撑日志汇聚场景。Serilog 没有非常官方的 Stream Sink但自己写一个并不难无非是ILogEventSink.Emit里StreamAddAsync把日志 JSON 塞进 Stream。注意 Redis 是单线程模型日志量超过每秒几十万条时吞吐会受限但中小规模完全够用。2.4 老系统中的 MSMQ 怎么办用自定义 Sink 兜底Windows 消息队列MSMQ虽然在 .NET Core 时代逐渐边缘化但确实还有不少老系统在用。Serilog 没有官方 MSMQ Sink好在自定义 Sink 的门槛极低。你需要一个类实现ILogEventSink在Emit里把日志渲染成消息写入 MessageQueue再写一个日志配置扩展方法把 Sink 挂进LoggerConfiguration。消息队列这个动作本身太简单了难点全在接入时的序列化、队列事务和异常处理。public class MsmqSink : ILogEventSink { private readonly MessageQueue _queue; private readonly IFormatProvider _formatProvider; public MsmqSink(MessageQueue queue, IFormatProvider formatProvider) { _queue queue; _formatProvider formatProvider; } public void Emit(LogEvent logEvent) { var text logEvent.RenderMessage(_formatProvider); var json JsonSerializer.Serialize(new { logEvent.Timestamp, logEvent.Level, Message text, Properties logEvent.Properties }); _queue.Send(json, MessageQueueTransactionType.Single); } }注意.NET 8环境下 System.Messaging 在非 Windows 系统不可用即使 Windows 上也需要额外的兼容包支持而且 MSMQ 服务本身需要开启。我对新项目的建议很明确新系统别选 MSMQ历史系统迁移期可以写一个自定义 Sink 做桥接让老队列里的日志继续流入新平台等窗口期过了再下线。这种自定义模式的思路也可以扩展到自研的消息系统任何有客户端库的队列都可以用同样方式接入 Serilog。为了让你选型时有更直观的参考我把常见方案的核心差异整理如下方案单机吞吐量参考运维复杂度消息可靠性适用规模RabbitMQ Sink中等万级/秒中等可配置持久化与 ACK中小型微服务集群Kafka Sink高十万级/秒高分区副本机制强日日志量亿级以上Redis Stream 自建 Sink中低万级以内低可持久化但受单线程限制中小规模日志汇聚Seq中低官方服务端保障日志量可控的单体/小集群MSMQ 自定义 Sink低中依赖 Windows 环境老系统迁移过渡3. .NET 8 接入 RabbitMQ Sink 的完整过程从 NuGet 到路由键设计3.1 最小接入一行配置把日志送进交换机先看最精简的接入方式。新建一个 .NET 8 的 ASP.NET Core 项目NuGet 里至少需要安装Serilog.AspNetCore和Serilog.Sinks.RabbitMQ如果后续要输出紧凑 JSON再加一个Serilog.Formatting.Compact。我们直接在Program.cs里配置using Serilog; var builder WebApplication.CreateBuilder(args); builder.Host.UseSerilog((context, services, configuration) { configuration .ReadFrom.Configuration(context.Configuration) .ReadFrom.Services(services) .Enrich.FromLogContext() .Enrich.WithMachineName() .WriteTo.RabbitMQ((clientConfig, sinkConfig) { clientConfig.Username logger; clientConfig.Password your-password; clientConfig.Exchange serilog.exchange; clientConfig.ExchangeType topic; clientConfig.RouteKey log; clientConfig.Port 5672; clientConfig.DeliveryMode DeliveryMode.NonDurable; sinkConfig.BatchPostingLimit 50; sinkConfig.Period TimeSpan.FromSeconds(2); }) .WriteTo.Console(); }); var app builder.Build(); app.Run();这段代码里UseSerilog接管了 ASP.NET Core 的日志管道所有ILoggerT的输出都会流经这套配置。.WriteTo.RabbitMQ后面那个 lambda 是分两个对象传参的一个是 RabbitMQ 客户端连接参数一个是 Sink 自身的批量参数别搞混。这里我保留了.WriteTo.Console()作为本地输出的兜底排查问题时能直接在控制台看到日志队列出问题也不会完全没有日志可看这是我在生产环境里一直保留的小习惯。3.2 appsettings.json 配置拆解BatchPostingLimit 与 Period 的含义如果你更习惯配置驱动上面的代码等价于在appsettings.json里这样配置{ Serilog: { Using: [ Serilog.Sinks.RabbitMQ ], MinimumLevel: { Default: Information, Override: { Microsoft.AspNetCore: Warning } }, WriteTo: [ { Name: RabbitMQ, Args: { exchange: serilog.exchange, exchangeType: topic, routeKey: log, username: logger, password: your-password, port: 5672, deliveryMode: Serilog.Sinks.RabbitMQ.DeliveryMode.NonDurable, batchPostingLimit: 50, period: 2 } } ] } }batchPostingLimit和period是批量 Sink 的灵魂参数。batchPostingLimit表示缓冲区里攒多少条日志后触发一次投递period表示最大等待时间。也就是说日志要么在攒满 50 条时立刻发出要么在 2 秒超时后把不满一批的日志也发出去。这两个参数决定了“吞吐”和“实时性”的平衡点。批量越大网络和队列交互次数越少吞吐越高但日志到达存储的延迟也越大。deliveryMode我强烈建议日志场景用NonDurable。原因很简单日志天生允许少量丢失尤其是 Debug 和 Information 级别的日志。RabbitMQ 的持久化要做到消息落盘写入性能会显著下降而日志系统最重要的指标是“扛得住峰值”不是“一条都不丢”。只有 Error 级或审计日志需要单独开持久化我们后面会提到用多 Sink 分流来区分处理。3.3 按日志级别分流交换机、路由键与队列绑定的实战设计路由键是 RabbitMQ 模型里最值得花心思设计的地方。一个常见的做法是用 topic 交换机路由键按log.level的格式组织比如log.information、log.warning、log.error。然后在 RabbitMQ 管理端创建三个队列分别绑定log.error.*、log.warning.*、log.information.*这样不同级别的日志自动进入到不同队列消费端可以按优先级处理。告警系统只消费 error 队列统计系统消费全量日志各取所需互不干扰。Serilog 侧可以通过多个 WriteTo 配合restrictedToMinimumLevel实现分级投递。还是用配置的方式举一个例子把 Error 日志单独投到一个 error 交换机其他日志走到默认交换机WriteTo: [ { Name: RabbitMQ, Args: { exchange: serilog.exchange, routeKey: log.error, deliveryMode: Serilog.Sinks.RabbitMQ.DeliveryMode.Durable, batchPostingLimit: 20, period: 1 }, restrictedToMinimumLevel: Error }, { Name: RabbitMQ, Args: { exchange: serilog.exchange, routeKey: log, deliveryMode: Serilog.Sinks.RabbitMQ.DeliveryMode.NonDurable, batchPostingLimit: 100, period: 2 } } ]Error 日志持久化且批次更小、投递更快普通日志不持久化用大批次换取吞吐。这个设计看起来简单实际效果极好线上告警的实时性和日志系统的吞吐量同时得到了保障。有一点务必记住Exchange、Queue、Binding 三件套最好在消费端或部署脚本里提前声明好Sink 只负责投递。3.4 结构化上下文Enricher 在消息队列场景下的角色日志进了消息队列之后最难的一步就是消费端不知道“这条日志是谁在什么上下文里打出来的”。Serilog 的 Enricher 就是为了解决这个问题。默认配置下我会挂这几个FromLogContext让LogContext里塞的属性跟着日志走WithMachineName记录物理机名WithThreadId记录线程WithSpan关联 .NET 的 Activity 追踪信息。这样消费端拿到 JSON 后能看到请求的 TraceId、SpanId、服务实例名等关键上下文排查跨服务问题时价值非常大。如果希望每条日志有一个全球唯一的业务标识可以自定义一个 Enricher这其实是一个很轻量的动作using Serilog.Core; using Serilog.Events; public class LogEventIdEnricher : ILogEventEnricher { public void Enrich(LogEvent logEvent, ILogEventPropertyFactory propertyFactory) { var eventId Guid.NewGuid().ToString(N); logEvent.AddPropertyIfAbsent( propertyFactory.CreateProperty(LogEventId, eventId)); } }注册方式也很简单在配置链路上加.Enrich.WithLogEventIdEnricher()。这个LogEventId字段在后面处理重复消费时是重要的幂等键。你可能会问Serilog 自带LogEventId不是有吗那个是事件源定义的短整型比如事件编号 1001不是全局唯一标识。排查问题、对接审计、去重时需要的是完全唯一的 ID自定义 Enricher 是最干净的做法。4. 消费端从队列到存储批量刷盘、字段规范化与背压4.1 一个能落地的消费 Worker批量缓冲与手动 ACK日志投递只是上半场消费端怎么接住才是真正的下半场。我的消费端是一个独立的 ASP.NET Core Worker 服务里面跑一个BackgroundService连接 RabbitMQ手动 ACK攒批落库。先看骨架代码public class LogConsumerWorker : BackgroundService { private readonly IConnection _connection; private readonly IChannel _channel; private readonly ListLogEntry _buffer new(); private readonly SemaphoreSlim _sync new(1, 1); protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _channel.BasicQos(0, 200, false); var consumer new AsyncEventingBasicConsumer(_channel); consumer.ReceivedAsync OnMessageReceived; _channel.BasicConsume(queue: serilog.queue, autoAck: false, consumer: consumer); while (!stoppingToken.IsCancellationRequested) { await Task.Delay(TimeSpan.FromSeconds(2), stoppingToken); await FlushBufferAsync(stoppingToken); } } private async Task OnMessageReceived(AsyncEventingBasicConsumer sender, BasicDeliverEventArgs args) { await _sync.WaitAsync(); try { var body Encoding.UTF8.GetString(args.Body.ToArray()); var entry JsonSerializer.DeserializeLogEntry(body); _buffer.Add(entry); if (_buffer.Count 500) { await FlushBufferAsync(CancellationToken.None); } } finally { _sync.Release(); } _channel.BasicAck(args.DeliveryTag, false); } private async Task FlushBufferAsync(CancellationToken token) { if (_buffer.Count 0) return; var batch _buffer.ToArray(); _buffer.Clear(); // 批量写入 Elasticsearch 或 ClickHouse拼接 Bulk 请求 await BulkWriteAsync(batch, token); } }BasicQos的第二个参数 200 是预取数量表示这个消费者一次最多从队列拉取 200 条未确认消息。这样做的目的是避免消费者内存里堆太多消息同时也让 RabbitMQ 在消费者异常退出时最多重投 200 条这个数字直接关系到下一节的重复消费数量。手动 ACK 放在批量写成功之后是关键点如果先 ACK 再写库写库失败消息就丢了如果先写库再 ACK写库成功但 ACK 丢失消息会被重复投递。日志链路上我一般选择“先写后 ACK”因为日志允许少量重复但不能接受大量丢失。4.2 日志 JSON 的规范化Properties 字典原来是这么一回事消费端拿到日志体之后第一件事是字段规范化。比如用CompactJsonFormatter序列化的日志时间戳通常是 ISO 8601 字符串级别是数字或字符串Properties是一个字典。如果直接把这个字典原封不动塞进数据库查询时会发现字段类型混乱同一个字段在一条日志里是字符串在另一条里是数组非常难处理。我的做法是消费端定义明确的LogEntry模型把t、l、m、x这些保留字段映射成强类型属性把Properties里的常用字段单独抽出来比如MachineName、ThreadId、TraceId、RequestPath。其余自定义属性统一放到一个Dictionarystring, object里落库时作为 map 字段存储。这样既保留了结构化日志的灵活性查询时又有稳定的主字段可用。如果你用的是 Elasticsearch可以在写入时配置 dynamic template 来约束字段类型如果用 ClickHouse建议直接把属性字段做成 Map 或者 JSON 列查询语法上需要提前规划好。4.3 背压处理消费端跟不上时宁可丢弃也不能阻塞业务消费端一定要考虑“队列积压”和“存储节点抖动”的情况。当 Elasticsearch 或数据库写入变慢时消费 Worker 内部的缓冲区会越涨越大最后内存爆掉。我的经验是提前给缓冲区设上限达到上限后采取降级策略而不是无限排队。降级策略有三个档次第一档是丢弃最旧的日志保证最新日志能入库第二档是把日志降级到本地文件等存储恢复后再补传第三档是触发告警让人来介入。实际项目中我用得最多的是第一档因为日志系统最怕的不是丢几条旧日志而是延迟过高导致问题发生时看不到实时日志。你要记住一个原则日志链路的任何环节都不允许反向阻塞业务线程如果队列积压已经超过可容忍范围丢日志比丢业务请求划算得多。4.4 消费端的重复消费与乱序提前给你打个预防针如果说背压是流量问题那重复消费就是可靠性问题。我遇到过一个很典型的案例消费端在处理一批日志时进程被 OOM Killer 杀掉这批日志已经写入了 Elasticsearch但 ACK 没有发出。进程重启后RabbitMQ 把那批消息重新投递于是同一批日志在 ES 里出现了两份。当时日志平台的对账功能弹了几十条重复告警排查了很久才定位到根因。重复消费在采用“至少一次投递”语义的消息队列里几乎无法完全避免。这也是为什么我要在日志链路里引入全局唯一的LogEventId。对于告警、审计这类对准确性要求较高的日志消费端必须做幂等处理对于普通日志可以根据业务容忍度决定要不要去重。下面的第五部分专门围绕这个问题展开。5. 消息队列重复消费问题成因、复现与日志场景的幂等方案5.1 三种消息队列各自制造重复的方式重复消费是消息队列的标配话题先说三种主流队列各自最容易制造重复的机制。RabbitMQ 在手动 ACK 模式下如果消费者在发送 ACK 之前断开连接或者 ACK 在网络传输中丢失未确认的消息会被 RabbitMQ 重新入队如果消费者进程在批量写库成功后、发送 ACK 前崩溃重复就必然发生。Kafka 的重复更多来自 offset 提交时机消费者处理完一批消息后如果还没来得及提交 offset 就发生再均衡重平衡后消费者组会从上次提交的 offset 位置重新消费这一段日志自然就被读了两遍。Redis Stream 的消费者组有类似机制消息被读取后会进入 PEL待处理消息列表只有显式 XACK 才会移除如果读取了消息但没确认重新消费时又会拿到。三者的共同本质是消息队列的投递语义通常是“至少一次”也就是保证不丢但不保证不重。日志链路要接受这个前提而不是寄希望于队列层帮忙去重。队列类型重复产生的典型时机触发条件重复粒度RabbitMQ消费端 ACK 前崩溃、连接断开未确认消息重新入队单条或一个预取批次Kafka消费组再均衡、offset 提交延迟从上次提交点重新拉取按分区的一段消息Redis StreamXACK 丢失、消费组重新读取PEL 中残留未确认消息单条消息5.2 实测复现手动 ACK 前杀掉进程会发生什么纸上谈兵不如亲自复现一次。我在测试环境搭了一套最小复现生产端用 Serilog 批量投递到 RabbitMQ消费端把预取值设为 500手动 ACK在批量写库之后、发送 ACK 之前直接杀掉进程。RabbitMQ 管理界面里能看到 500 条消息处于 unacked 状态进程消失后这些消息的状态从不确认变成 ready重新入队。等消费端再次启动它会再次收到这 500 条消息入库后日志总数多了 500 条。第二次复现我模拟的是网络抖动消费端在 ACK 发出瞬间断网ACK 没有到达 RabbitMQ。消息同样被重新投递但因为消费端已经重启且日志写入是幂等键保护下的重复才没有产生实际影响。这个实验给我一个非常直观的结论只要使用手动 ACK 的至少一次投递模型重复消费就不是“会不会发生”的问题而是“什么时候发生”的问题日志系统必须默认消息可能重复。5.3 日志去重要不要做基于 LogEventId 的幂等设计日志场景对去重的要求和其他业务不太一样。业务订单不能重复提交所以必须严格幂等日志少重复几条对检索和排查几乎无感。所以在日志链路上我建议分级处理Information 及以下级别的日志允许重复率在 0.1% 以内不做额外去重开销Warning 级日志可以适当做基于时间窗口的轻量去重Error 级日志和审计日志必须严格幂等因为这些日志会触发告警、参与对账重复消息会直接导致告警风暴。严格幂等的做法也不复杂。生产端给每条日志加上全局唯一的LogEventId消费端在写入存储前先去重。用一个伪代码来说明写入数据库时的处理逻辑public async Task WriteDeduplicatedAsync(LogEntry entry) { // Redis 60 秒短期去重 var key $log:{entry.LogEventId}; var added await _redis.StringSetAsync(key, 1, TimeSpan.FromSeconds(60)); if (!added) { // 说明 60 秒内已经写过了 return; } // MySQL / PostgreSQL 唯一索引兜底 await _db.ExecuteAsync( INSERT IGNORE INTO logs(log_event_id, level, message, timestamp) VALUES (EventId, Level, Message, Timestamp), entry); }Redis 的StringSet带过期时间做了第一道短期去重数据库的唯一索引log_event_id做了第二道长期兜底。INSERT IGNORE的意思是冲突时静默跳过不会因为重复日志抛异常导致消费进程终止。这个方案的成本在哪里每条日志多一次 Redis 操作和一次数据库写入对消费端吞吐有影响。所以我的实践是只对 Error 级日志走这个流程普通日志直接批量写入性能代价完全可控。5.4 用 Redis 与数据库唯一索引做双保险的写法如果你觉得两条链路都做太重至少要做到“数据库唯一索引 UPSERT”。这在 MySQL、PostgreSQL、ClickHouse 里都有对应语法。核心思想是让存储层对LogEventId唯一重复数据插入时要么忽略、要么覆盖。日志场景建议忽略因为我们要的是原始第一条覆盖可能把上下文信息冲掉。还有一种场景值得注意同一个LogEventId被重复投递时第一次写入可能是在旧索引里第二次写入时索引配置已经变了。这种跨索引的重复去重不能只靠唯一索引必须在写入前通过目标索引的 document ID 做查询判断。Elasticsearch 的_id可以直接指定为LogEventId重复写入时使用create操作而不是index操作create遇到已存在文档会返回冲突错误而不会覆盖。这也是我在日志平台中处理重复问题的首选方案因为它不需要额外的 Redis 组件完全依赖 ES 自身的乐观锁语义。6. 压测一轮之后吞吐量、关键参数与两个隐性成本6.1 压测方法与指标口径入队速率、消费速率、落库延迟配置和代码都就位后一定要压测不然根本不知道系统真实的承载边界在哪。我的压测方式是在一个独立环境里跑一个控制台程序循环生成符合业务分布的日志用Log.Information和Log.Error按 95:5 的比例混合抛给 Serilog持续 30 分钟。监控指标主要看三个RabbitMQ 管理界面的 publish rate 和 queue depth、消费端每秒入库条数、从日志产生到落库的端到端延迟。在我常用的 4 核 8G 三节点测试环境里RabbitMQ Sink 开启 50 条批量后单实例投递能稳定在 8000 到 12000 条每秒消费端批量 500 条写入 Elasticsearch入库速度大约 6000 条每秒。这个数据的意义不是让你直接套用而是帮你建立量级概念批量 Sink 以每秒万条为基准线如果你的日志峰值超过这个数要么加实例用多个 Sink 分片要么考虑换 Kafka。6.2 RabbitMQ 侧参数调优prefetch、队列持久化与 TTLRabbitMQ 侧的参数直接影响稳定性和峰值吸收能力。第一个是队列的持久化策略日志队列不建议开持久化前面反复说过了但如果业务要求日志尽量不丢可以折中为“交换机持久化 队列持久化 消息非持久化”这样至少队列和交换机在 broker 重启后还在消息丢了但架构不会乱。第二个是队列长度上限一定要设置x-max-length或x-message-ttl让积压日志在超过一定量或一定时间后自动丢弃避免 MQ 变成第二个被拖垮的磁盘。第三个是消费者预取数量我刚才用的 200 只是一个起点如果单条日志体积较大预取数量要调小防止消费端内存被撑爆。还有一个容易被忽略的参数是 RabbitMQ 的心跳超时。默认心跳 60 秒如果日志量大或网络不稳长时间没有活动时连接可能被服务端判定超时并断开。建议把心跳超时调大比如 120 秒同时启用自动重连。Serilog 的 RabbitMQ Sink 一般通过底层的 ConnectionFactory 参数控制压测时一定要把连接重连的场景验证一遍否则一次网络抖动就能让日志静默丢失。6.3 消费端参数调优线程数、批量阈值与缓冲区上限消费端的参数要比生产端更敏感。批量阈值不能只看每秒条数还要看单条日志的体积。我们的日志平均体积大约 1KB500 条一批就是 500KB写入 Elasticsearch 的 Bulk 请求正好在合理区间但如果日志体积到了 10KB批大小得降到 100 条左右防止内存峰值过高。消费线程数也不是越多越好我测过 2 个线程和 8 个线程的差异发现存储端写入速度才是瓶颈线程数超过 4 以后收益迅速衰减反而增加了上下文切换和连接占用。缓冲区上限在这里是一个安全阀。线上我配置的是 Channel 容量 10000 条满员时不再从队列拉消息而不是在内存里无限堆积。存储恢复后消费端自动从队列继续读取积压由 RabbitMQ 队列兜底。这个设计保证消费端进程本身不会因为日志积压而 OOM队列积压则是可观测、可弃、可扩容的比进程崩溃要容易处理得多。6.4 容易被忽视的序列化开销与 GC 抖动压测过程中我发现两个隐性成本特别值得警惕。第一个是日志模板渲染开销。RenderMessage会把{Request}、{UserId}这些占位符渲染成字符串如果模板里塞了很大的对象或者集合每次记录日志都会产生一次代价高昂的ToString()调用。我见过有人把整个请求体对象直接塞进日志模板结果单条日志渲染消耗了几十毫秒业务线程被日志拖慢。正确做法是记录关键字段或者用前缀让 Serilog 做结构化存储而不是直接渲染字符串。第二个是 GC 压力。批量日志场景每秒产生数万个字符串对象Gen0 和 Gen1 回收非常频繁极端情况下日志吞吐会直接反映到 GC 暂停里。我在生产环境把 Server GC 打开后日志写入对应用暂停的影响有明显缓解。此外尽量复用JsonSerializer实例避免每条日志都创建新的序列化器。还有一点RabbitMQ Sink 自身的批处理机制虽然提高了吞吐但也让日志在内存中滞留的时间变长批量值越大停顿期间内存占用越高batchPostingLimit和period的搭配需要根据实例内存反复调。写到这里消息队列 Sink 的链路基本完整了。根据我在生产环境跑下来的体会日志应该走消息队列的项目往往都有一个共性日志吞吐量已经超过单机存储的稳定写入能力。如果只是想收集日志做透视先别上 RabbitMQ文件加 Seq 就够用如果已经在为日志写库发愁那消息队列中间层基本是必经之路。最后再分享一个小技巧RabbitMQ 的日志队列一定要配上最大长度或 TTL日志千万不能让队列无限制堆积我们线上就是靠 queue TTL 加死信队列把过期日志自动丢弃才没让积压日志反过来成为另一场事故。
阅读完成 · 觉得有帮助?
咨询建站