KEFK
KEFK 通常指 Kafka、Elasticsearch、Fluentd 和 Kibana。名称只列出核心组件,实际部署还需要边缘采集器;常见链路是 Filebeat 或 Fluent Bit 采集,Kafka 缓冲,Fluentd 聚合处理后写入 Elasticsearch,Kibana 查询展示。
架构
Fluentd 也可以同时部署为边缘 Agent 和中心 Aggregator,但资源敏感的容器环境通常使用 Fluent Bit 作为 DaemonSet,中心端再使用 Fluentd 做复杂路由和插件处理。若团队已经统一使用 Filebeat,也可以直接输出 Kafka,不必为了名称完整而更换采集器。
Kafka 带来的变化
本篇关注 Kafka 在日志平台中的角色。Kafka 自身的 Partition、Consumer Group、Offset、副本和投递语义,参见消息队列板块的 Kafka。
| 维度 | 直接写 Logstash/Fluentd | 经过 Kafka |
|---|---|---|
| 架构复杂度 | 较低 | 增加 Broker、Topic、Partition 和消费组运维 |
| 削峰 | 依赖 Agent 与处理器本地队列 | 可吸收更长时间和更大规模积压 |
| 回放 | 本地队列通常不适合任意回放 | 保留期内可重置 Offset 重放 |
| 多消费者 | 需要额外复制或多路输出 | 不同消费组独立消费同一日志 |
| 故障隔离 | Elasticsearch 故障较快传到采集端 | 消费端可暂停,生产端继续写到容量边界 |
| 成本 | 组件少 | 机器、存储、网络和治理成本更高 |
日志量不大、只有一个消费端且 Logstash Persistent Queue 足够覆盖故障窗口时,没有必要仅为“高可用”引入 Kafka。存在明显流量尖峰、需要多下游消费、要求小时级回放,或 Elasticsearch 维护窗口较长时,Kafka 的价值才明显。
多主机不等于多下游
Filebeat 的同一个 Output 可以配置多个 hosts,用于负载均衡和故障切换。例如多个 Logstash 地址仍属于同一条输出链路,事件不会各复制一份供不同系统独立处理。
一个 Filebeat 实例通常只启用一种 Output。需要让同一份日志同时进入检索、归档、实时计算和安全分析时,不应在每台主机上堆叠多个采集进程来复制数据,而应先写入 Kafka,再用不同 Consumer Group 独立消费:
同一 Consumer Group 内的多个实例用于分摊分区,不会让每个实例都收到完整副本;不同 Consumer Group 才会各自读取完整日志流,并维护独立 Offset。这是 Kafka 实现一对多分发的关键。
需要进一步判断确认级别、重复消费、幂等和故障恢复边界时,参见 Kafka 可靠性。
Topic 与分区设计
- Topic 按数据域、保留等级或安全边界划分,例如
logs.application、logs.audit,不要为每个实例创建 Topic。 - 消息 Key 可使用服务名或稳定业务维度;Key 过于集中会产生热点分区。
- 分区数决定消费并行度,但过多分区会增加 Broker、客户端和控制面的开销。
- Retention 必须大于可接受的下游最长恢复时间,并按压缩后实际写入量估算磁盘。
- 应用日志通常允许至少一次投递,消费者必须容忍重复;审计事件需要更严格的 ID、保留和校验策略。
Filebeat 写入 Kafka
filebeat.inputs:
- type: filestream
id: application-logs
paths: ["/var/log/apps/*.json"]
parsers:
- ndjson:
target: ""
add_error_key: true
queue.disk:
max_size: 2GB
output.kafka:
hosts: ["kafka-1:9093", "kafka-2:9093", "kafka-3:9093"]
topic: "logs.application"
key: "%{[service.name]}"
required_acks: -1
compression: gzip
max_retries: 10
ssl.enabled: true
ssl.certificate_authorities: ["/etc/filebeat/certs/ca.crt"]生产端确认成功只表示 Kafka 已按确认策略接收消息,不表示 Elasticsearch 已建立索引。必须分别监控生产失败、Broker 副本状态、消费组 Lag、消费失败和最终检索延迟。
Fluentd 消费与写入
下面展示中心 Fluentd 的关键结构。Kafka 与 Elasticsearch 插件需要随镜像显式安装并锁定版本:
<source>
@type kafka_group
brokers kafka-1:9092,kafka-2:9092,kafka-3:9092
topics logs.application
consumer_group log-indexer
format json
tag logs.application
</source>
<filter logs.application>
@type record_transformer
<record>
data_stream.type logs
data_stream.dataset application
data_stream.namespace prod
</record>
</filter>
<match logs.application>
@type elasticsearch
hosts https://es-1:9200,https://es-2:9200
user fluentd_writer
password "#{ENV['ELASTICSEARCH_PASSWORD']}"
data_stream_enable true
<buffer topic,partition>
@type file
path /var/lib/fluent/buffer/elasticsearch
flush_thread_count 4
flush_interval 5s
retry_type exponential_backoff
retry_forever true
</buffer>
</match>具体参数取决于所用 Fluentd 插件版本,部署前应以镜像中插件的配置文档为准。密码通过 Secret 注入,不写入配置仓库。
Fluentd 的 Buffer
Fluentd Buffer 把输入与输出解耦,优先使用持久化文件缓冲而不是纯内存缓冲。需要明确:
- Chunk Key 决定事件如何分块,过细会制造大量小文件。
- Buffer 上限决定能够承受多长时间的 Elasticsearch 故障。
- Retry 超限后的行为必须明确,不能静默丢弃。
- Overflow 策略会影响采集端阻塞或丢弃,应与 Kafka 保留能力联合设计。
Kafka 已承担主缓冲时,Fluentd Buffer 主要保护正在处理的批次和短时输出抖动,不需要无限放大;消费 Offset 应在消息安全进入输出链路后再提交。
失败与回放
回放前应先修复 Mapping、解析或权限问题,再使用新的消费组验证小范围数据。直接把旧消费组 Offset 大幅回退可能造成重复写入、集群突发压力和告警风暴。
关键监控
| 层次 | 指标 |
|---|---|
| 采集器 | 活跃 Harvester、读取延迟、队列占用、发送失败 |
| Kafka | Under-replicated Partition、磁盘、生产延迟、消费组 Lag |
| Fluentd | Buffer 使用率、Retry、异常 Chunk、处理速率 |
| Elasticsearch | 写入延迟、Rejected、分片、磁盘水位、生命周期失败 |
| 端到端 | 日志产生时间到可检索时间、抽样完整率、重复率 |
只监控组件存活不足以判断日志平台可用。应周期性写入带唯一 ID 的探针日志,并验证它能在规定时间内被 Kibana 或 Elasticsearch 查询到。