日志系统看似不起眼,但一旦业务流量上来,它往往是最先暴露问题的地方:采集端写入阻塞、日志服务被打挂、磁盘 IO 被写满。把 Kafka 引入日志链路,用它做采集端的缓冲层,是业界非常成熟的做法。这篇文章会结合实际项目经验,讲清楚 Kafka 在日志采集场景下的架构设计、关键配置和常见坑。

一、为什么日志链路需要 Kafka 做缓冲
很多团队的最初方案是应用直接把日志写到本地文件,再用采集工具(比如 Filebeat、Fluentd)转发到集中式存储。这种方式在流量平稳时没什么问题,但有两个隐患:一是下游存储(如 ES、ClickHouse)写入速度有限,流量高峰时采集端会被反压甚至阻塞;二是下游一旦故障,日志就在采集端堆积甚至丢失。
Kafka 的价值在于把“生产”和“消费”彻底解耦。应用或采集器只管往 Kafka 里写,写入成功即代表日志已落盘持久化(多副本机制保证可靠性),下游消费者按自己的节奏消费。哪怕 ES 挂了半小时,日志也只是安静地待在 Kafka 里,等恢复后继续消费,一条不丢。
一个典型的日志链路是:应用产生日志 -> Filebeat/Fluentd 采集 -> Kafka -> Logstash/Flink 消费 -> ES/ClickHouse/对象存储。Kafka 在中间扮演的角色就是削峰填谷的蓄水池,这也是它区别于直接写存储方案的核心优势。
二、生产端配置:批量、压缩与可靠性
日志场景的特点是量大、单条价值相对低、允许极小概率丢失(视业务而定)。因此生产端的配置思路是:优先保证吞吐,同时通过批量发送摊薄单条开销。下面是一份经过生产验证的 Filebeat 输出到 Kafka 的配置:
filebeat.inputs:
- type: log
enabled: true
paths:
- /var/log/app/*.log
output.kafka:
hosts: ["kafka1:9092", "kafka2:9092", "kafka3:9092"]
topic: "app-log"
required_acks: 1
compression: gzip
max_message_bytes: 1000000
batch_size: 16384
channel_buffer_size: 4096
几个参数值得展开说明。compression 设置为 gzip 后,日志文本通常能压缩到原始大小的 20% 左右,能显著降低网络和磁盘压力;如果 CPU 资源紧张,可以改用 snappy 或 lz4,压缩率略低但 CPU 开销更小。required_acks 设为 1 表示 leader 副本写入成功即返回,吞吐和可靠性之间取了折中;如果日志绝对不能丢,建议改成 all,配合 topic 的 min.insync.replicas 使用。
如果是应用直接通过客户端发送,Java 客户端的核心配置类似,重点同样是批量:
Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// 批量配置:一批最多 64KB,最多等 50ms 凑批
props.put("batch.size", 65536);
props.put("linger.ms", 50);
// 压缩
props.put("compression.type", "lz4");
// 缓冲区总量,按吞吐量调整
props.put("buffer.memory", 67108864);
// 发送失败重试
props.put("retries", 3);
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
linger.ms 是新手最容易忽略的参数。默认为 0,即立即发送,每条消息一次网络往返,吞吐很低。设为 20 到 50 毫秒后,生产者会稍微等待以凑满批次,单条延迟只增加几十毫秒,但吞吐往往能提升数倍。日志场景对这几十毫秒完全不敏感,强烈建议开启。
三、Topic 分区设计与消费端负载均衡
分区数直接决定并行度的上限。一个经验公式是:分区数取“消费端期望的总吞吐除以单个消费者的吞吐”再向上取整,并且至少等于消费者实例数,否则多余实例会空闲。日志场景通常按小时几 GB 的量来估算,设 12 到 32 个分区比较常见。
创建 topic 时还要考虑两个参数:retention.ms 控制日志保留时间(比如保留 3 天),retention.bytes 控制按大小清理,两者取先到者。日志 topic 一般不需要紧凑型存储,但副本因子建议至少 2,避免单机故障导致数据不可用。
kafka-topics.sh --create \ --bootstrap-server kafka1:9092 \ --topic app-log \ --partitions 24 \ --replication-factor 2 \ --config retention.ms=259200000 \ --config retention.bytes=107374182400
消费端用消费者组实现负载均衡:同一个 group.id 的多个实例,每个分区只会被其中一个实例消费。这里有个常见的坑——如果日志需要同时写入 ES 和对象存储,千万别在同一个消费组里做双写,一旦某一路写入变慢就会拖慢整体消费。正确做法是起两个不同 group.id 的消费者组,各自独立消费,互不影响。
Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092");
props.put("group.id", "log-to-es");
props.put("enable.auto.commit", "false");
props.put("max.poll.records", 500);
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props,
new StringDeserializer(), new StringDeserializer());
consumer.subscribe(List.of("app-log"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> record : records) {
// 写入 ES 或其他存储
}
consumer.commitSync();
}
四、消息积压排查与容量规划
上线后最常见的问题是消费滞后,也就是消费速度跟不上生产速度。排查分两步:先看监控确认滞后量,kafka-consumer-groups.sh --describe 能直接看到每个分区的 CURRENT-OFFSET、LOG-END-OFFSET 和 LAG,配合 Prometheus + kafka_exporter 可以做可视化告警。
确认积压后要定位原因。如果消费端 CPU 没跑满,多半是下游写入慢,比如 ES 索引刷新间隔太频繁、bulk 批次太小,把 refresh_interval 调大到 30 秒、bulk 每批凑够 5 到 15MB 通常能明显改善。如果消费端 CPU 已经跑满但单个分区消费能力到了上限,就需要扩容消费者实例或者增加分区数——注意增加分区会破坏按 key 的有序性,日志场景一般不依赖顺序,影响不大。
# 查看消费组各分区滞后情况 kafka-consumer-groups.sh \ --bootstrap-server kafka1:9092 \ --describe \ --group log-to-es
容量规划上给出几个经验数字作为起点:3 节点集群、普通 SSD 磁盘,单 broker 顺序写吞吐可达每秒几百 MB,对多数业务来说日志流量远够。真正需要提前算的是磁盘总量:保留 3 天、每天 200GB 原始日志、压缩后约 50GB,那么集群总磁盘需求约 300GB,再预留 30% 的余量给突发流量和 rebalance 期间的日志膨胀。
五、几个实战踩坑总结
第一,消息体过大。日志里偶尔会出现超长堆栈或大 JSON,超过 message.max.bytes(默认约 1MB)会被 broker 拒绝,采集端表现为静默丢弃。建议在生产端截断超长日志,或者调大 broker 和 topic 两级的字节限制(topic 级不能超过 broker 级)。
第二,频繁 rebalance。消费者处理太慢导致两次 poll 间隔超过 max.poll.interval.ms(默认 5 分钟),消费者会被踢出组触发 rebalance,形成恶性循环。解决办法是减小 max.poll.records,或者把慢的处理逻辑改成异步批量提交。
第三,副本滞留导致 ISR 收缩。如果某个 broker 磁盘 IO 持续跟不上,follower 副本会被移出同步副本列表,当 min.insync.replicas 无法满足时生产端会直接报错。监控里一定要加上 ISR 数量告警,提前发现磁盘瓶颈。
整体来说,Kafka 做日志缓冲的方案本身并不复杂,关键在于参数调优贴合业务流量特征,并且把滞后量、ISR 健康度、磁盘使用率这几项监控配齐。搭好之后,这套管道在面对十倍流量突增时依然能稳稳接住,这也是它被大规模普及的根本原因。