拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Filebeat+Kafka+ClickHouse:PB级日志平台实战总结

Filebeat+Kafka+ClickHouse:PB级日志平台实战总结

做淘客返利APP最怕什么?流量高峰一来,订单日志却查不动。年初我们线上出过一次事故:佣金结算数据整整延迟了40分钟,客服电话被打爆,技术群里每秒钟都有人在刷屏。最后定位到根因,是老的日志采集分析方案在峰值下直接崩溃——Elasticsearch集群的写入把CPU打满,查询一条订单日志要等十几秒,数据积压越来越多,最终把整个链路拖死。

那次事故之后,我下定决心把日志平台彻底重构。最终落地的方案就是标题里那套组合:Filebeat + Kafka + ClickHouse,目标是支撑PB级日志数据的实时采集、缓冲和秒级检索。这套架构上线后跑了快一年,再也没因为日志系统出过大事故。这篇文章就是我把整个过程——从选型、架构设计、部署踩坑到查询优化、容量运维——完整梳理出来的一份实战总结,希望能给正在做日志平台、或者准备把业务日志从ES迁移到ClickHouse的朋友一些参考。

1. 事故复盘:日志系统为什么会成为业务瓶颈

1.1 淘客返利APP的日志到底有什么特殊性

先说说业务背景。淘客返利APP的核心链路是:用户通过推广链接下单 -> 平台收到订单回推 -> 校验佣金 -> 结算返利。这中间每一步都在产生日志,大致可以分成这么几类:

  • 埋点日志:用户的点击、浏览、分享、唤起APP行为,量大、字段固定;
  • 订单回推日志:电商平台通过开放接口推送订单状态,包含订单号、商品、金额、推广位等信息;
  • 支付与佣金日志:支付回调、佣金比例计算、结算状态流转,业务价值最高,也最需要实时查询;
  • API访问日志:网关层记录所有请求的耗时、状态码、来源IP。

这些日志有一个共同特点:平时看着平稳,一到活动节点就瞬间爆发。比如晚上8点到10点的黄金时段、大促返利加码活动,峰值QPS能冲到平时的5到8倍。订单类日志又是强实时业务——用户刚下单就希望看到返利进度,运营要实时监控佣金异常,客服要秒查用户订单情况。所以日志平台不光是"能存",还得"能查得快"。

数据量方面,我们的设计目标是:日均日志事件量百亿级,高峰期每秒几十万条,离线存储按PB级别规划。这个量级下,任何"灵活有余、性能不足"的存储方案都会被轻松击穿。

1.2 旧方案到底崩在哪

事故之前,我们的日志链路是典型的ELK变种:Filebeat -> Redis -> Logstash -> Elasticsearch 5.x。Filebeat采集日志写到Redis,Logstash从Redis消费再写入ES,查询走Kibana。

这套方案在小流量下没什么问题,但到了大促场景就暴露出一堆毛病:

  • ES写入吞吐成为瓶颈:日志场景写入量大、单个文档又小(几百字节),ES的索引和分词开销极高,CPU一会儿就打满,写入延迟伴随副本同步延迟一起飙升;
  • Redis做缓冲太脆弱:Logstash消费速度跟不上时,Redis内存不断上涨,最后触发LRU淘汰,日志直接丢;
  • 查询能力过剩但性能不足:ES最擅长的是全文检索和复杂条件聚合,但我们90%的查询都是"按订单号精确查""按用户ID查""按时间统计"这种固定模式,根本用不上分词和全文索引,反而被索引开销拖累;
  • 磁盘和内存成本太高:ES副本机制加上倒排索引,实际磁盘占用是原始数据的好几倍,堆内存动不动就是几十GB,运维压力非常大。

说白了,不是ES不好,是我们在错误的场景里用了它。日志查寻这种"海量写入 + 固定模式查询 + 冷热数据分层"的需求,应该交给更专一的引擎。

2. 技术选型:为什么偏偏是Filebeat + Kafka + ClickHouse

2.1 采集端:Filebeat凭什么赢

采集端当时有三个候选:Logstash、Fluentd、Filebeat。我们的选择标准很务实:轻量、配置简单、稳定性高。

Logstash能做复杂的加工清洗,但它是JVM进程,单机内存占用轻松上G。集群规模一旦上来,光采集端的资源消耗就非常可观。Fluentd生态不错,但Ruby环境的部署和插件管理在运维上不如Go二进制的Filebeat省心。Filebeat采用的是Go编写、单二进制运行,内存占用通常在几十MB到一两百MB,CPU几乎可以忽略。配置是纯YAML,采集什么文件、输出到哪、加什么字段、正则多行合并,全在配置文件里搞定,改完reload就行。

另外Filebeat自带backpressure机制:当输出端(Kafka)写不进去的时候,它会自动放缓读取速度,而不是拼命把日志读进内存——这一点在生产环境非常关键,能有效防止雪崩。

2.2 缓冲层:Kafka在消息队列选型中的位置

日志链路里的消息队列,考察重点是吞吐量、可靠性、消费生态。市面上主流MQ我基本都对比过:

维度KafkaRocketMQRabbitMQ
吞吐量单节点几万到几十万条/秒高,但部署较重几千到几万条/秒,更适合业务解耦
消费模型Consumer Group + offset管理Consumer GroupQueue + Exchange机制更灵活
日志场景适配顺序写盘、批量消费、长期保存更偏业务消息、事务消息适合小规模、复杂路由
运维成本依赖ZooKeeper(或KRaft),集群管理成熟需要NameServer和Broker双组件相对简单

日志场景本质上是"高吞吐的流式数据搬运",Kafka的PageCache利用、顺序写盘、批量拉取这些设计,简直是为日志量身定做的。它的offset机制让消费进度可以精确管理,挂了也能续跑不丢数据。RocketMQ在金融级事务消息上更强,RabbitMQ在业务解耦、消息路由上更灵活,但论"把海量日志稳稳地从A搬到B",Kafka是最省心的选择。

2.3 存储引擎:ClickHouse vs Elasticsearch vs Doris

存储层的对比才是这次重构的核心。我把三者做了个横向对比:

维度ClickHouseElasticsearchDoris
存储模型列式存储倒排索引 + 文档存储列式存储(MPP)
压缩比极高,日志场景常见压缩5-10倍较低,副本和索引放大明显较高
写入吞吐非常高,批量写入轻松上百万行/秒受限于索引和分片,易成瓶颈高,但部署架构复杂
查询模式固定模式聚合查询极快(秒级)全文检索和模糊搜索强适合多表Join的OLAP分析
运维复杂度单机即可起步,集群靠ZooKeeper协调节点多、内存大、运维复杂组件多(FE/BE)、部署门槛高
日志场景匹配度极高中低中

ClickHouse最打动我的几点,第一是列式存储带来的超高压缩比——日志里有大量重复的URL、状态码、用户ID,列存压缩之后,同样的原始日志,磁盘占用大约是ES的1/5到1/10;第二是MergeTree家族天然适合时间序列数据,按时间分区、异步合并、TTL过期删除,几乎不需要额外开发;第三是查询是真正的秒级,比如统计某个推广位一天产生的订单量和佣金总额,SQL一条搞定,200亿行数据也就一两秒。

Doris在复杂Join分析上更强,但日志场景几乎没有Join需求,加上部署维护成本更高,我就放弃了。

3. 核心链路设计:数据从采集到可查询的每一步

3.1 数据流总览

整个链路我用一句话描述:文件采集 -> 消息缓冲 -> 批量写入 -> 列式存储 -> SQL查询。

具体展开是:

  • 应用服务器上跑着Filebeat,采集本地日志文件,做简单的多行合并、字段清洗;
  • Filebeat把数据以Kafka Producer的身份写入Kafka指定Topic,天然削峰填谷;
  • 自研的Go消费端从Kafka拉取数据,批量写入ClickHouse;
  • ClickHouse按天分区存储,通过TTL管理生命周期;
  • 上层业务通过一个检索API统一查询ClickHouse,返回订单追踪、用户行为、佣金统计等结果。

这套设计里每个环节的职责都单一、清晰,任何一个组件挂了都能独立恢复,数据不会因为单点故障丢得干干净净。

3.2 Filebeat采集配置的实战要点

先贴一份我们生产环境的Filebeat配置骨架:

filebeat.inputs: - type: log enabled: true paths: - /data/logs/app/order-service/*.log fields: app: order-service env: prod fields_under_root: true multiline.pattern: '^\[20\d{2}-\d{2}-\d{2}' multiline.negate: true multiline.match: after output.kafka: hosts: - "10.x.x.21:9092" - "10.x.x.22:9092" - "10.x.x.23:9092" topic: "app_log_ingest" partition_hash: hash: ["fields.order_no"] compression: gzip max_message_bytes: 1000000 required_acks: 1

几个容易踩坑的细节:

  • multiline一定要配。Java服务端日志里的异常堆栈常常跨多行,如果不做多行合并,一条异常会被拆成十几条垃圾数据,下游解析时全部报错。合并正则基于时间戳开头判断即可。
  • fields_under_root要慎重。把app、env提到根上,后续在ClickHouse里可以当作字段直接过滤,不用解析嵌套JSON。
  • registry文件不能乱删。Filebeat通过/var/lib/filebeat/registry记录每个文件读到了哪个offset,重启不会重读。如果误删,会导致全量重新采集,Kafka和下游直接被打爆。

3.3 Kafka的Topic和分区策略设计

Kafka侧我们用了三个Topic分别承载不同类型日志,避免业务互相挤兑:

Topic承载数据分区数副本数保留时间
app_log_ingestAPI、系统访问日志12248小时
app_order_log订单回推、支付回调24272小时
app_bury_log埋点点击曝光24224小时

分区数不是越大越好。分区越多,消费者并行度越高,但Broker上的文件句柄和ISR同步开销也越大。我们按"高峰期每秒各Topic写入量 / 单分区可承载的5万条每秒"再乘上安全余量来定,控制在12~24个分区。

Kafka分区策略上,同一个订单号的消息必须进同一个分区。因为消费端要保证同一订单的日志按时间顺序处理,如果你乱分到不同分区,Data Race和乱序问题会搞到你怀疑人生。我们通过partition_hash按order_no做哈希,保证同单同行。

另外生产环境建议开启compression.type=lz4,日志类消息压缩比高,能显著降低带宽和Broker磁盘压力。

3.4 消费写入层:为什么自研而不是直接用Kafka Engine

ClickHouse本身提供Kafka Engine,可以直接从Kafka消费写表。但我强烈不建议在日志这种高吞吐场景直接用它,原因有三:

  • 难以控制批量大小:Kafka Engine消费节奏比较任性,无法按业务低峰高峰灵活调整写入批次,容易产生大量小parts,直接导致ClickHouse的Merge压力暴涨;
  • offset管理不透明:Kafka Engine的offset是保存在集群内的,一旦遇到数据格式变化、表重建,offset丢了或者重复消费都很难排查;
  • 错误处理能力弱:如果某条数据格式有问题,Kafka Engine默认会阻塞Topic,影响整个消费链路。

所以我们自研了一个Go写的消费端,核心逻辑不复杂:

func main() { // 从Kafka拉取消息 reader := kafka.NewReader(kafka.ReaderConfig{ Brokers: []string{"10.x.x.21:9092", "10.x.x.22:9092", "10.x.x.23:9092"}, GroupID: "clickhouse-writer", Topic: "app_order_log", MinBytes: 1e4, // 10KB MaxBytes: 10e6, // 10MB }) batch := make([]LogRecord, 0, 5000) for { msg, err := reader.ReadMessage(context.Background()) if err != nil { continue } record, parseErr := parseMessage(msg.Value) if parseErr != nil { // 打点记录坏消息, 不阻塞消费 continue } batch = append(batch, record) // 攒够5000条或超过3秒, 才真正落库 if len(batch) >= 5000 || time.Since(lastFlush) > 3*time.Second { clickhouse.BatchInsert(batch) batch = batch[:0] } } }

批量写入是核心。ClickHouse在批量写入(每次几万行)时吞吐极高,但如果你一行一行写,系统会非常难受,还可能触发Too Many Parts错误。我们把单批次控制在5000~20000行之间,结合async_insert、max_insert_threads等参数,实测写入速度轻松跑满单节点几百万行每秒。

消费端的幂等也要考虑。Kafka的at-least-once语义决定了极端情况下会有重复消息,好在我们的日志场景天然可以接受"少量重复"——查询日志重复一两条不影响业务判断。

3.5 ClickHouse表结构设计:分区、排序键与TTL

以订单日志为例,这是我们最终的建表语句:

CREATE TABLE app_order_log ( order_no String, app_id UInt32, event_date Date, event_time DateTime, user_id UInt64, product_id UInt64, amount Decimal(18,4), status String, trace_id String, raw_message String ) ENGINE = MergeTree() PARTITION BY event_date ORDER BY (event_time, order_no) TTL event_time + INTERVAL 420 DAY;

几个关键设计决策:

  • PARTITION BY event_date:按天分区,查询天然带时间过滤时,可以直接跳过无关分区。删除过期数据时,DROP PARTITION的成本远低于DELETE WHERE。
  • ORDER BY (event_time, order_no):排序键决定主索引稀疏索引的粒度。日志查询几乎都带时间范围,把event_time放第一顺位是最优解。
  • TTL直接写在表上:保留420天,ClickHouse后台会自动把过期数据删掉,不需要额外写定时任务。
  • raw_message保留原始数据:有时候结构化字段解析不全,保留原始日志用于兜底排查,代价是多一些磁盘占用,值得。

4. 部署实操:三件套安装与踩过的那些坑

4.1 ClickHouse 21.8的部署与几个典型坑

我们用的是ClickHouse 21.8.15.7,这个版本列式存储性能稳定,部署方式也很成熟。安装没什么好说的,rpm包装完改下config.xml和users.xml就能起来。真正麻烦的是运行参数,当时踩了几个典型的坑:

第一个坑:Too Many Parts。上线初期我们写入频率过高、每次批量又太小,ClickHouse后台MergeThreads来不及合并,导致/var/lib/clickhouse/data/下某个表分区堆积了几千个parts,直接触发too many parts保护,写入开始报错。解决办法是调大单次写入批次、降低写入频率,同时把merge_with_ttl_timeout适当调大,让后台合并更积极。

第二个坑:ZooKeeper抖动导致插入慢。ClickHouse的分布式表依赖ZooKeeper做元数据和协调,21.8时代ZK抖动会直接影响insert和alter。我们单独部署了三节点ZooKeeper,并把ZK的堆内存、会话超时参数做了加固,好很多。

第三个坑:count(*)很慢。日志场景下大家习惯性想查总行数,但MergeTree的count(*)需要扫描数据。我们在业务侧改成查SELECT sum(rows) FROM system.parts WHERE table='app_order_log',秒出,避免每次大促看总量时把集群查卡。

4.2 Kafka集群安装与调优细节

Kafka我们用的2.8版本,三节点起步。部署上的一个经验:Kafka更适合把数据目录放到独立的高速磁盘上,SSD最好,机械盘会直接拉垮吞吐。另外JVM堆内存不需要给太大,Kafka大量使用PageCache,堆内存设8~16G就够,多了反而浪费。

生产环境重点配置:

log.retention.hours=48 num.io.threads=8 num.network.threads=8 default.replication.factor=2 min.insync.replicas=1 message.max.bytes=1000000 log.segment.bytes=1073741824

日志场景我们的acks设置为1,因为允许少量丢失,换取写入吞吐。如果是订单、佣金这类关键链路,建议相关Topic单独设置acks=all+min.insync.replicas=2,但吞吐会有一定牺牲,需要按业务场景权衡。

4.3 数据接入层的联调注意点

整个链路联调时,最容易忽视的是消息格式的前后一致性。我建议从第一天起就在Kafka消息里统一生成trace_id和event_time,Filebeat采集时在文件里就有,消费端解析顺便就取了,下游查询时统一按这个字段做业务追踪,省掉后续大量格式兼容的麻烦。

5. 检索平台的SQL优化:PB级数据下如何做到秒级响应

5.1 查询场景决定了优化方向

日志检索平台的查询不像BI报表那样五花八门,我们的需求90%都集中在这几类:

  • 按订单号查完整流转记录(用户下过单吗、支付成功没、佣金计算了没);
  • 按用户ID和时间范围查某段时间的行为轨迹;
  • 按推广位/活动ID维度统计订单数、金额、佣金;
  • 查某条trace_id关联的整条调用链路。

场景明确之后,优化就有靶子了:精确查询走索引,统计查询走物化视图。

5.2 跳数索引:让精确查询变成真正的秒级

日志表的数据量上来之后,即使有分区裁剪,单个分区内可能还有几亿行数据。这时候ORDER BY里的稀疏索引只能帮你快速定位到某个时间范围,没法定位到具体订单号。我们的做法是给高频查询字段加跳数索引:

ALTER TABLE app_order_log ADD INDEX idx_order_no (order_no) TYPE bloom_filter GRANULARITY 4; ALTER TABLE app_order_log ADD INDEX idx_user_id (user_id) TYPE bloom_filter GRANULARITY 8;

Bloom Filter索引的原理是"快速判断这个分区块内有没有可能包含目标值",有则扫描,没有则跳过。加了这套索引之后,按订单号精确查询的响应时间从几十秒降到了几百毫秒。注意GRANULARITY参数别设太小,否则索引文件本身会占用大量空间,通常4或8比较合理。

5.3 物化视图:统计报表的必备加速器

实时统计类查询,比如"今天每个推广位产生多少订单、多少佣金",直接在原始表上GROUP BY是很大开销的,每次扫全量数据重算一遍毫无必要。

我们建了一套物化视图,底层用AggregatingMergeTree做预聚合:

CREATE MATERIALIZED VIEW mv_order_daily ENGINE = AggregatingMergeTree() PARTITION BY event_date ORDER BY (event_date, app_id) AS SELECT event_date, app_id, count() AS order_cnt, sum(amount) AS total_amount, sumIf(amount, status = 'paid') AS paid_amount FROM app_order_log GROUP BY event_date, app_id;

插入原始表时,数据会同步进入物化视图的聚合状态,之后查统计就是查预聚合结果,千万行变几百行,查询自然就是毫秒级。

注意:ClickHouse的物化视图更像"插入触发器",它不会自动更新旧数据。所以务必保证物化视图和原始表同时创建,不要在原始表已经有历史数据之后再补建。

5.4 慢查询排查的基本套路

遇到慢查询,我的排查顺序是固定的:

EXPLAIN SELECT ... FROM app_order_log WHERE ...;

先看ReadFromMergeTree这一步的parts数量和扫描行数。如果扫描行数异常大,优先检查分区裁剪是否生效、排序键是否覆盖了过滤字段、索引是否命中。另外,日志表的low_cardinality字段(比如status、app_id)建议声明成LowCardinality(String),能显著提升压缩率和过滤速度。

6. 稳定性和容量规划:PB级日志系统的日常运维心法

6.1 积压监控:防止Kafka把下游冲垮

日志平台的稳定性核心在Kafka积压量。消费端如果跟不上生产速度,Kafka的lag会持续增长,最后积压几千万条,ClickHouse被一次性灌入大量数据,parts直接爆炸。

我设计了两个监控指标:Kafka Topic的Lag(消费落后)和消费者的每秒消费速率。用Prometheus + Grafana把消费者Lag暴露出来,超过阈值就告警。同时消费进程设置背压:当Kafka Lag减少到接近0时,自动降低消费速率,避免无限追尾。

6.2 磁盘容量怎么估算

容量规划是个实操问题,计算公式其实不复杂:

  • 单条日志原始大小:约500字节(JSON格式);
  • 峰值每秒写入量:50万条,即每秒250MB原始数据;
  • 一天原始数据量:约21TB(按峰值算,实际全天均值按峰值的30%算,约6TB/天);
  • ClickHouse列存压缩比按1:5算,实际一天落库约1.2TB;
  • 保留420天,总存储约500TB,考虑多副本和合并膨胀余量,预留550TB。

我们初期按这个模型规划了三组机器,之后每月核对一次实际压缩比和日增数据量,动态扩容。

6.3 parts合并和写入之间的平衡

ClickHouse的数据是先在内存中攒成一个小part(默认100~1000行),然后异步合并成大part。写入越碎,parts越多,合并压力越大。我们日常运维会关注system.parts中活跃parts数量,当单分区parts超过300就开始干预:降低插入频率、加大批量、或者临时暂停一些小查询释放合并线程。前期的容量规划里给磁盘留足余量也很有用,因为clickhouse在merge时对磁盘临时空间的需求大概是数据量的1.5倍。

6.4 权限和数据生命周期管理的沉淀

日志数据在合规视角下需要明确生命周期。我们的做法是ClickHouse里每个库按部门/业务线隔离,用户账号只授予必需的库表只读或写入权限;原始日志按保留期自动TTL清理,审计类、佣金结算类日志单独设置更长保留期。数据治理的理念如果到后期才补,会非常痛苦,建议一开始就在表结构设计时留好app_id、env这类隔离字段。

7. 沉淀下来的经验与我们的下一步打算

整套系统从上线到现在,最直观的感受是:技术选型要对得起业务场景。日志平台最核心的诉求是海量写入、压缩存储、固定查询,ClickHouse在这些维度上几乎是为这个场景量身定做的。Filebeat + Kafka的组合则保证了整个链条任何时候都不会因为流量波动而雪崩,层次清晰,每一层出现问题都能独立恢复。

再分享几个只有自己动手踩坑才懂得的体会:

  • 日志平台的数据格式越早统一越好。我们在项目初期就定了Json编码、统一字段名、统一时间格式,这省掉了后期无限多的兼容性工作;
  • 不要把Kafka Engine当成生产级日志接入方案。它适合做轻量数据管道和演示验证,生产环境一定要有可控的批量写入层;
  • 物化视图的聚合逻辑要跟业务一起评审。统计口径错了,返工成本很高,宁可前期多花时间梳理指标定义;
  • 容量和带宽的评估一定要在高峰期做。平常看着很宽裕的资源,在大促峰值面前可能就是薄薄一层窗户纸。

目前这套架构还有一个我们正在迭代的方向:引入轻量的数据质量校验,在消费端对日志进行实时规则校验,脏数据打标隔离而不阻塞主链路;另外把埋点类日志逐步拆成独立Topic,未来接入实时数仓做用户路径分析。这套骨架的可扩展性足够,后续演进的空间还很大,也希望有类似场景的朋友能少走些弯路。

返回列表