Skip to content

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.applicationlogs.audit,不要为每个实例创建 Topic。
  • 消息 Key 可使用服务名或稳定业务维度;Key 过于集中会产生热点分区。
  • 分区数决定消费并行度,但过多分区会增加 Broker、客户端和控制面的开销。
  • Retention 必须大于可接受的下游最长恢复时间,并按压缩后实际写入量估算磁盘。
  • 应用日志通常允许至少一次投递,消费者必须容忍重复;审计事件需要更严格的 ID、保留和校验策略。

Filebeat 写入 Kafka

yaml
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 插件需要随镜像显式安装并锁定版本:

text
<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、读取延迟、队列占用、发送失败
KafkaUnder-replicated Partition、磁盘、生产延迟、消费组 Lag
FluentdBuffer 使用率、Retry、异常 Chunk、处理速率
Elasticsearch写入延迟、Rejected、分片、磁盘水位、生命周期失败
端到端日志产生时间到可检索时间、抽样完整率、重复率

只监控组件存活不足以判断日志平台可用。应周期性写入带唯一 ID 的探针日志,并验证它能在规定时间内被 Kibana 或 Elasticsearch 查询到。