设计一个日志收集系统
提出问题
"设计一个日志收集系统"是系统设计面试里最容易被低估的题目。看起来简单——不就是把日志从各个服务拉到中央存起来么?但面试官问的是:百万 QPS 写入怎么扛?ELK 和 Loki 怎么选?日志存储成本爆炸了怎么办?全链路追踪怎么打通? 这些问题背后本质是成本 × 性能 × 可观测性的三元平衡。Uber 在 2020 年公开过他们的日志系统每天处理 500 TB 日志,用了三层降采样和冷热分离才把成本压住。小公司无脑上 ELK 的结果往往是 ES 集群每月几万块的存储费,实际上 90% 的日志从来没人查过。
分析问题
核心链路:采集 → 传输 → 缓冲 → 存储 → 检索
日志收集系统的主干就这五步,每一步都有对应的成熟组件:
采集层 → 传输层 → 缓冲层 → 处理层 → 存储层 → 检索层
│ │ │ │ │ │
Filebeat Kafka Kafka Logstash ES/Loki Kibana/Grafana
Fluentd Kafka (分区) (过滤) (冷热) (查询)采集层:在每台机器上部署轻量级 Agent(Filebeat / Fluentd),只做一件事——读日志文件、发到 Kafka。Agent 的内存占用要控制在 50MB 以内,不能影响业务进程。Filebeat 的 harvester 模式用指针跟踪文件偏移量,崩溃重启后从断点续传,不丢日志。
# filebeat.yml 精简配置
filebeat.inputs:
- type: log
paths:
- /var/log/app/*.log
fields:
service: user-service
env: production
output.kafka:
hosts: ["kafka-1:9092", "kafka-2:9092"]
topic: "app-logs"
partition.round_robin:
reachable_only: true传输与缓冲层:Kafka 是日志系统的核心瓶颈点。百万 QPS 写入,Kafka 用分区数扛——每个分区顺序写盘,吞吐随分区数线性增长。一个 3 节点 6 分区的 Kafka 集群轻松扛 50 万条/秒,瓶颈通常在网卡带宽。关键参数:acks=1(不等待所有副本确认,日志丢几条可以接受,但性能提升 10 倍)、compression.type=snappy(压缩比 3:1,减少网络带宽)。
处理层:Logstash 从 Kafka 消费日志,做解析、过滤、脱敏、路由。这里有个常见的坑:Logstash 的 Ruby 处理引擎性能很差,300MB 堆内存每秒只能处理 2 万条。生产上用 Vector(Timber 出品,Rust 写的)替代 Logstash,吞吐 10 倍提升,内存只有 1/10。
# vector.toml 配置示例
[sources.kafka]
type = "kafka"
group_id = "log-consumer"
topics = ["app-logs"]
[transforms.parse]
type = "remap"
inputs = ["kafka"]
source = '''
. = parse_json!(.message) ?? {}
.timestamp = now()
.level = del(.level) ?? "INFO"
'''
[transforms.redact]
type = "remap"
inputs = ["parse"]
source = '''
.message = replace!(.message, r'1[3-9]\d{9}', "***手机号***")
'''
[sinks.elasticsearch]
type = "elasticsearch"
inputs = ["redact"]
endpoint = "http://es-hot:9200"
index = "logs-%Y.%m.%d"ELK vs Loki 选型
这是面试官必然会问的对比题,不会只答"ELK 全文搜索、Loki 轻量级"就够。
| 维度 | ELK (Elasticsearch) | Loki (Grafana Loki) |
|---|---|---|
| 索引方式 | 全文倒排索引,每个字段都建索引 | 只存元数据 label,日志内容不建索引 |
| 存储成本 | 高(1TB 原始日志→ES 约 2-3TB,含副本) | 低(1TB 原始日志→Loki 约 300-500GB) |
| 查询能力 | 全文搜索、聚合、DSL 灵活强大 | 按 label 检索 + LogQL 正则,不支持全文搜索 |
| 写入吞吐 | 受限于索引写入,需要调优 refresh_interval | 纯 append,吞吐很高 |
| 运维复杂度 | 高(集群调优、分片管理、GC 调优) | 低(无状态组件,水平扩展简单) |
| 典型场景 | 需要全文搜索异常栈、复杂分析查询 | 按服务/标签查原始日志、告警关联 |
选型决策树:
需要全文搜索(查异常栈、查关键词)?
├─ 是 → ELK(选型完成)
└─ 否 → 需要复杂聚合分析?
├─ 是 → ELK
└─ 否 → 只需要按标签+时间查原始日志?
├─ 是 → Loki
└─ 否 → ELK(保险方案)很多公司最终走的是 ELK + Loki 双通道:ERROR 级别日志双写两份(一份进 ES 用于搜索,一份进 Loki 用于低成本归档),INFO/DEBUG 日志只写 Loki。这样搜索能力保留,存储成本降低 60%。
日志降采样
百万 QPS 下全量存储成本无法承受,降采样是必选项。面试官想听的不是"按级别采样"这种基础方案,而是业务感知采样。
哈希采样(TraceId 一致性):同一个 TraceId 的日志要么全存要么全丢,保证单个请求的日志上下文完整性。用 hash(TraceId) % 100 < 采样率 判断。
// 哈希采样实现
public class HashSampler {
private final int sampleRate; // 1-100
private final String traceId;
public boolean shouldSample() {
int hash = Math.abs(traceId.hashCode()) % 100;
return hash < sampleRate;
}
}动态采样(按级别分层):
- ERROR:100% 全量存储
- WARN:10% 采样
- INFO:1% 采样
- DEBUG:默认不采集,通过动态配置开关按需开启
关键路径采样:只存核心接口(用户下单、支付等)的完整日志,下游服务(推荐、广告等)只存异常日志。这个方案依赖全链路 TraceId 的精准标注,在网关层给核心请求打上 critical=true 标签,采集侧根据标签决定采样率。
# 动态采样配置(可从配置中心热更新)
SAMPLE_RULES = {
"default": { # 默认采样率
"ERROR": 1.0,
"WARN": 0.1,
"INFO": 0.01,
"DEBUG": 0.0
},
"critical": { # 核心链路采样率
"ERROR": 1.0,
"WARN": 1.0,
"INFO": 1.0,
"DEBUG": 0.1
}
}冷热分离
日志存储的"二八定律":80% 的查询落在最近 7 天的日志上,但 30 天以前的数据占了 70% 的存储空间。
ES ILM(Index Lifecycle Management)实现冷热分离:
hot(SSD,7 天)→ warm(HDD,30 天)→ cold(S3/MinIO,90 天)→ delete// ES ILM 策略
{
"policy": {
"phases": {
"hot": {
"min_age": "0ms",
"actions": {
"rollover": { "max_size": "50GB", "max_age": "1d" },
"set_priority": { "priority": 100 }
}
},
"warm": {
"min_age": "7d",
"actions": {
"allocate": { "require": { "box_type": "warm" }},
"set_priority": { "priority": 50 }
}
},
"cold": {
"min_age": "30d",
"actions": {
"allocate": { "require": { "box_type": "cold" }},
"set_priority": { "priority": 0 },
"freeze": {}
}
},
"delete": {
"min_age": "90d",
"actions": { "delete": {} }
}
}
}
}冷热分离不是简单的"SSD 和 HDD 的区别"。冷数据节点可以大幅降低副本数(热数据 2 副本,冷数据 1 副本甚至 0 副本),因为冷数据丢了可以重算。如果日志是写到 S3 的,S3 本身有 11 个 9 的持久性,所以 ES 冷阶段副本都可以不存,直接从 S3 查 S3 快照。
TraceId 全链路透传
日志收集的终极价值是排障。没有 TraceId 的日志就是一盘散沙,出了线上问题要去 ES 里搜关键字,全凭运气。
透传方案:网关层生成 TraceId → 通过 HTTP Header 传递 → 各服务通过 MDC(Mapped Diagnostic Context)注入到日志线程上下文。
// Spring Boot 过滤器自动注入 TraceId
@Component
public class TraceIdFilter implements Filter {
@Override
public void doFilter(ServletRequest req, ServletResponse res, FilterChain chain) {
HttpServletRequest request = (HttpServletRequest) req;
String traceId = request.getHeader("X-Trace-Id");
if (traceId == null || traceId.isEmpty()) {
traceId = UUID.randomUUID().toString().replace("-", "");
}
MDC.put("traceId", traceId);
try {
// 调用下游服务时透传 TraceId
chain.doFilter(new TraceIdRequestWrapper(request, traceId), res);
} finally {
MDC.remove("traceId");
}
}
}日志采集时,Filebeat 从日志的 traceId 字段提取,写入 ES 的 trace_id 字段。查询时直接 trace_id: "xxx" 就能拿到一个请求经过的所有服务的日志,按时间线排序,排障效率提升 10 倍。
PII 脱敏
日志中可能包含手机号、身份证、银行卡号等敏感信息,直接存 ES 不脱敏等于裸奔。
脱敏方案分三层:
- 采集侧:Filebeat 采集前用正则替换(性能好,但规则维护成本高)
- 处理侧:Logstash/Vector 用
grok或remap做脱敏(推荐,统一管理规则) - 查询侧:ES 的
search-template返回时脱敏(兜底,但数据已经落地了)
// 采集侧脱敏 - Logstash filter 配置
filter {
mutate {
gsub => [
"message", "(1[3-9]\\d{9})", "***手机号***",
"message", "(\\d{6}[\\dXx]{4}\\d{4}[\\dXx]{4}\\d{4}[\\dXx])", "***身份证***"
]
}
}总结
日志收集系统设计的关键在于理解成本与能力的平衡:
- 采集层用 Filebeat,轻量不侵占业务资源
- 传输层用 Kafka 做缓冲,分区数扛高吞吐
- 存储层 ELK vs Loki 按场景选,双通道是性价比最优解
- 降采样用哈希 + 动态 + 关键路径三层组合,不丢关键日志
- 冷热分离用 ILM 自动迁移,冷数据存 S3 比 HDD 还便宜 10 倍
- TraceId 透传是排障能力的基础,不是可选项
- PII 脱敏在 Logstash/Vector 处理层做,统一管理
生产上最常见的翻车场景:日志量预估不足。上线前说日均 1TB,上线后业务量暴涨到 10TB,Kafka 分区不够、ES 索引写入跟不上、磁盘被打满。提前规划好降采样策略和冷热分离方案,而不是等磁盘告警了再补。
参考:ELK 官方架构指南、Loki 设计与定位文档、Uber Logging 实践、ES ILM 冷热分离最佳实践