Filebeat+Kafka+ClickHouse:PB级日志平台实战总结
发布时间:2026/10/2 19:46:13来源:尧图网络
做淘客返利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采集日志写到RedisLogstash从Redis消费再写入ES查询走Kibana。这套方案在小流量下没什么问题但到了大促场景就暴露出一堆毛病ES写入吞吐成为瓶颈日志场景写入量大、单个文档又小几百字节ES的索引和分词开销极高CPU一会儿就打满写入延迟伴随副本同步延迟一起飙升Redis做缓冲太脆弱Logstash消费速度跟不上时Redis内存不断上涨最后触发LRU淘汰日志直接丢查询能力过剩但性能不足ES最擅长的是全文检索和复杂条件聚合但我们90%的查询都是按订单号精确查按用户ID查按时间统计这种固定模式根本用不上分词和全文索引反而被索引开销拖累磁盘和内存成本太高ES副本机制加上倒排索引实际磁盘占用是原始数据的好几倍堆内存动不动就是几十GB运维压力非常大。说白了不是ES不好是我们在错误的场景里用了它。日志查寻这种海量写入 固定模式查询 冷热数据分层的需求应该交给更专一的引擎。2. 技术选型为什么偏偏是Filebeat Kafka ClickHouse2.1 采集端Filebeat凭什么赢采集端当时有三个候选Logstash、Fluentd、Filebeat。我们的选择标准很务实轻量、配置简单、稳定性高。Logstash能做复杂的加工清洗但它是JVM进程单机内存占用轻松上G。集群规模一旦上来光采集端的资源消耗就非常可观。Fluentd生态不错但Ruby环境的部署和插件管理在运维上不如Go二进制的Filebeat省心。Filebeat采用的是Go编写、单二进制运行内存占用通常在几十MB到一两百MBCPU几乎可以忽略。配置是纯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搬到BKafka是最省心的选择。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拉取数据批量写入ClickHouseClickHouse按天分区存储通过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.typelz4日志类消息压缩比高能显著降低带宽和Broker磁盘压力。3.4 消费写入层为什么自研而不是直接用Kafka EngineClickHouse本身提供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 tableapp_order_log秒出避免每次大促看总量时把集群查卡。4.2 Kafka集群安装与调优细节Kafka我们用的2.8版本三节点起步。部署上的一个经验Kafka更适合把数据目录放到独立的高速磁盘上SSD最好机械盘会直接拉垮吞吐。另外JVM堆内存不需要给太大Kafka大量使用PageCache堆内存设8~16G就够多了反而浪费。生产环境重点配置log.retention.hours48 num.io.threads8 num.network.threads8 default.replication.factor2 min.insync.replicas1 message.max.bytes1000000 log.segment.bytes1073741824日志场景我们的acks设置为1因为允许少量丢失换取写入吞吐。如果是订单、佣金这类关键链路建议相关Topic单独设置acksallmin.insync.replicas2但吞吐会有一定牺牲需要按业务场景权衡。4.3 数据接入层的联调注意点整个链路联调时最容易忽视的是消息格式的前后一致性。我建议从第一天起就在Kafka消息里统一生成trace_id和event_timeFilebeat采集时在文件里就有消费端解析顺便就取了下游查询时统一按这个字段做业务追踪省掉后续大量格式兼容的麻烦。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未来接入实时数仓做用户路径分析。这套骨架的可扩展性足够后续演进的空间还很大也希望有类似场景的朋友能少走些弯路。
网站建设高端定制企业官网