直接用 go func() { send(msg) }() 会丢数据,因 goroutine 无持久化、无 ack、无重试,进程崩溃后任务静默消失;必须先落盘或写入 kafka/redis stream,再由消费者处理并手动确认。

为什么直接用 go func() { send(msg) }() 会丢数据
不是协程启动慢,是进程挂了就全丢。goroutine 本质是内存里的执行单元,没持久化、没 ACK、没重试——机器重启、OOM、panic 后,正在跑的几十个 send() 全部静默消失,下游根本收不到。你看到“发出去了”,只是把任务塞进了 goroutine 调度队列,不是真送达。
真实场景里,订单通知、短信推送、埋点上报这些必须“至少一次”送达的任务,不能靠 go 关键字兜底。
- 必须把任务先落盘或写入 Redis Stream / Kafka 等中间件,再由独立消费者拉取
- 消费者拿到消息后,要先做业务处理,成功后再
XACK或CommitOffset - 失败则不 ACK,让中间件重投;超时未 ACK 的,由
XCLAIM或死信队列接管
batchProcessor 函数里最容易漏掉的计时器重置
看标准实现:用 time.NewTicker(timeout) 监听超时,但每次 out 后,如果没手动重置 ticker,下一次超时仍从上轮开始计时——比如设了 5 秒超时,第 3 秒就凑满 batch 发走了,那剩余 2 秒的倒计时还在跑,2 秒后又触发一次空发(<code>len(batch) == 0),或者更糟:下一批第一条消息进来后,离上次发送才过 1 秒,却要等满 4 秒才发。
正确做法是在每次批量发送后调用 ticker.Reset(timeout),且只在因数量触发发送时重置(超时触发本身已是新周期起点)。
- 数量触发发送 → 必须
ticker.Reset(timeout) - 超时触发发送 → 不用重置,ticker 自动续期
- 输入通道关闭且
len(batch) > 0→ 发完即 return,不重置
批量大小和超时时间不是拍脑袋定的
设 batchSize = 100、timeout = 5 * time.Second 看似合理,但实际取决于下游吞吐能力。比如调用一个平均耗时 80ms 的 HTTP 接口,100 条并发发过去,很可能触发对方限流 429;而设成 20 条 + 200ms,QPS 更稳,P99 延迟反而更低。
关键是要监控两个指标:batch_size_actual(实际发出批次的平均大小)和 batch_latency_ms(从第一条进缓存到整批发出的延迟)。如果前者长期
- 压测时用固定 QPS 模拟上游,观察下游错误率和延迟拐点
- 生产环境把 batchSize 和 timeout 做成可热更新的配置项,别硬编码
- 对延迟敏感的业务(如实时风控),优先保 timeout,batchSize 可动态收缩(例如当前批次 300ms 内只攒到 12 条,就发这 12 条)
别在批量发送路径里做日志格式化和文件 I/O
批量处理的核心价值是减少系统调用次数,但很多人在 out 前加一行 <code>log.Info("send batch", "size", len(batch), "first_id", batch[0].ID),这就把异步批处理退化成了同步打点——每次都要 time.Now().Format()、拼字符串、锁 log buffer,甚至触发 GC。
真正高吞吐的推送服务,日志只记关键信号:batch_sent、batch_dropped、batch_retry,且字段全用预计算整数(如 unixnano)、避免任何字符串操作。
- 业务日志走另一条异步通道,和推送逻辑完全隔离
- 错误日志(如 HTTP 503、连接 refused)必须绕过批量通道直写,否则失败时连“为什么失败”都看不到
- 所有磁盘 I/O(包括轮转日志)必须异步,
lumberjack默认的同步压缩会卡住整个 writer goroutine
批量处理的难点从来不在“怎么攒数据”,而在“怎么不破坏上下游 SLA”——超时控制错位、日志拖慢主路径、下游扛不住还硬推,这些问题比 channel 怎么写更致命。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











