1. 项目背景:淘客返利APP日志平台的选型思路
先说下这个项目的背景。我做过的这个淘客返利APP,用户日常操作包括浏览商品、点击推广链接、下单、确认收货、返利到账等一系列行为。整个链路里,日志数据的体量非常夸张——光是一天的用户点击流日志就能到几十亿条,加上订单状态变更、返利结算记录、系统运行日志,数据量很快就冲到了PB级。
在这个量级下,日志平台要解决的已经不是"把日志存下来"这么简单,而是三个核心诉求:一是实时性,用户刚操作完,运营马上就能看到数据;二是查询效率,客服排查用户订单纠纷、运营分析渠道转化,都得秒级出结果;三是存储成本,PB级数据如果选错存储引擎,光磁盘开销就能把成本拖垮。
结合这几个诉求,我最终定下的方案是Filebeat采集日志,Kafka做消息缓冲和解耦,ClickHouse做存储和检索。这套链路在最开始设计的时候参考了日志数据量、流式处理需求和成本预算三个维度。淘客返利这种业务形态的特殊性在于:闲时流量和活动大促期间的峰值流量差距能到十倍以上,日志链路必须能扛住这种脉冲式的流量冲击。Kafka在这条链路里承担的不只是消息管道,更是一个天然的流量缓冲池,让下游ClickHouse写入压力平滑可控。
1.1 为什么不用Logstash或Flume
在采集端的选型上,团队里有人提议用Logstash,也有人提Flume,但最终都否掉了。
Logstash在数据清洗和格式化上确实强,但它是JVM系的应用,默认堆内存动不动就分配2到4个G。我们的采集Agent要部署到每台业务服务器上,如果每台机器上都跑一个吃2G内存的进程,那对业务应用的资源侵占太明显了。淘客返利APP的业务服务器配置普遍是16G到32G内存,Logstash占比太高,业务侧肯定不干。相比之下,Filebeat是Go语言实现的,常驻内存能控制在30M到50M左右,几乎可以忽略不计。这就是一个"杀鸡用牛刀还是用手术刀"的问题。
Flume的问题在于它本身就是为大数据生态设计的,配置文件繁琐,source、channel、sink三段式配置写起来一套一套的,部署运维都偏重。而对于"读文件、发Kafka"这个需求,Filebeat一行配置就能搞定,社区活跃度也更高。
1.2 为什么用ClickHouse而不是Elasticsearch
存储和检索引擎的选型,团队内部也做过多轮充分论证对比。
Elasticsearch在日志检索领域确实是老牌选手,但到了PB级这个体量,它的成本问题非常突出。ES的倒排索引天生就要额外占用大量存储,一份原始日志存进去,加上索引和副本,物理磁盘占用往往是原始数据的3倍左右。ClickHouse的列式存储结构天然适合压缩,LZ4压缩算法下日志类数据的压缩比能做到5比1到10比1,玄学一点说,同样一块1T的磁盘,ES能存300G的日志,ClickHouse能存500G以上,这就是本质的成本差异。
再从检索场景来看,ES擅长的全文检索在我们的业务里用得并不多。淘客返利日志的查询更像是结构化查询——查某个用户ID在某时间段的下单记录,查某个渠道的PV、UV、订单转化率,这些都是典型的列式聚合查询,是ClickHouse的主场。ClickHouse的聚合分析性能在百万到亿级数据量下能到毫秒到秒级,完全能满足运营看板和客服排查的需求。
另外,ClickHouse还有一个Elasticsearch没法比的硬核优势——数据写入吞吐量。我们在压测环境验证过,单机ClickHouse的写入吞吐能做到每秒5万到10万行,配合批量写入还能更高。而ES在高并发写入下分片容易产生热点,频繁触发refresh反而让写入性能下降,量大的时候还得靠调bulk size和线程数慢慢调优。
1.3 整条链路的架构定位
这套方案在整体架构里的定位可以理解为:Filebeat是毛细血管,负责把各处的"血液"收集起来;Kafka是动脉,负责高速运输和缓冲;ClickHouse是心脏,负责数据的高效落地和查询服务。链路总共就三段,中间没有XML配置、没有JVM调优、没有复杂的拓扑管理,出了问题排查路径很短。
在真实的淘客返利场景里,这种简洁的链路价值非常大。有一次大促活动,某个接口的日志量瞬间翻了五倍,如果是传统架构,日志管道早就被冲垮了。而我们的链路里,Kafka的partition数量是按照峰值流量预留的,Filebeat本地还有内存缓冲和磁盘缓冲兜底,ClickHouse的写入线程池也留了余量。那次大促整个链路一个告警都没报,稳得一批。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. Filebeat采集端配置与实操要点
Filebeat虽然轻量,但要把日志完整、不丢、不重地送到Kafka,配置上是有几道坎要过的。下面我把实际配置一步步拆开讲。
2.1 filebeat.yml核心配置详解
yaml复制filebeat.inputs:
- type: log
enabled: true
paths:
- /data/applogs/*.log
fields:
app_name: taoke_app
log_type: biz_order
fields_under_root: true
encoding: utf-8
ignore_older: 12h
scan_frequency: 10s
output.kafka:
hosts: ["kafka1:9092", "kafka2:9092", "kafka3:9092"]
topic: "app_log_topic"
partition.hash:
hash: ["app_name"]
compression: lz4
max_message_bytes: 1048576
required_acks: 1
worker: 2
bulk_max_size: 2048
logging.level: info
logging.to_files: true
logging.files:
path: /var/log/filebeat
name: filebeat.log
keepfiles: 5
几个关键参数说明一下:
fields_under_root: true 这个参数容易忽略,但很关键。它决定自定义字段是放在JSON的顶层还是嵌套在fields子对象里。如果不设置,查询的时候就要写成 fields.app_name,设置之后直接查 app_name 就行,在ClickHouse建表时字段映射会省很多事。
scan_frequency 是扫描新文件的频率。需求是对日志近实时采集,设置成10秒,就是在新增日志到被发现的延迟是10秒,业务的容忍范围内。如果set成1s,会带来磁盘IO开销,没必要。
bulk_max_size 控制单次发送到Kafka的消息条数,默认是2048。这个值不是越大越好,太大了会增加内存缓冲的压力,太小了会频繁创建网络连接,吞吐上不去。我之前遇到过调到8192后Filebeat进程频繁GC和内存溢出的问题,后来回到2048才稳定。
worker 是每个输出Kafka的并发worker数。默认是1,我调成2是为了让单机Filebeat的发送吞吐更高,但前提是目标Kafka能接得住。如果下游消费扛不住,调大worker只会加快堆积,没有任何意义。
2.2 多行日志合并的难点
淘客返利APP的日志里有很多异常堆栈,占用多行。如果每个堆栈行都当成独立事件发给Kafka,后续在ClickHouse里查询时,一个异常就变成了几十行数据,根本没法分析。
yaml复制multiline:
pattern: '^\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}'
negate: true
match: after
这个意思是以"时间戳开头"为一条新日志的标志,如果某一行不是以时间戳开头,就归属到上一条日志的尾部。配置起来很简单,但它背后的逻辑要想明白:如果业务日志里有非标准前缀的普通行,这个pattern就要调整。我遇到过一种情况,日志里有一行是连续的"---"分割线,这个分割线被合并到了上一条日志的尾部,导致日志内容里出现一块没意义的字符。后来通过调整正则把分割线也排除掉才解决。
2.3 Filebeat调优和资源控制
Filebeat的资源消耗虽然很低,但在高吞吐场景下还是要留心。首先是文件句柄数,每扫描一个文件就要占用一个fd,如果日志文件数量巨大,要调大系统的ulimit -n。其次是内存缓冲,Filebeat内部有一个默认的mem queue,容量是4096条消息,如果日志量持续大于这个数,会产生背压,采集会开始阻塞。这时候优先思路是提高bulk_max_size,让每次发送的消息更多,而不是盲目加内存。
还有一个小细节是close_inactive。如果一个日志文件5分钟没有新内容,Filebeat默认会关闭它的文件句柄。这本身没什么问题,但当一个文件在关闭之后又重新有内容写入时,Filebeat可能因为状态记录里的偏移量问题导致漏采。为了保险,我在核心业务日志的配置里把close_inactive设置成了1h,代价是fd占用稍微高一点,但避免漏日志的风险。
3. Kafka在链路中的角色与集群配置
Kafka在整个链路里是规则的中枢,它既承担了缓冲压力、削峰填谷的职责,也决定了数据的最终分区归属。这部分配置不好,后面的ClickHouse查询性能分分钟被拖累。
3.1 Kafka集群部署参数经验
淘客返利日志场景,我给Kafka集群规划的是3节点,配置是8核16G内存,每节点挂3块SAS盘做raid5。磁盘这块单独说一下——Kafka对延迟最敏感的就是磁盘IO,SSD虽然贵,但如果你对实时性有高要求,建议至少上SSD。用机械盘的话,峰值写入时段的磁盘IO很容易到瓶颈,整个集群的吞吐会掉一个档次。
Kafka的server.properties里有几个参数在日志场景下值得关注:
code复制num.partitions=12
default.replication.factor=2
min.insync.replicas=1
log.retention.hours=24
log.segment.bytes=1073741824
num.partitions 默认是1,日志场景下太少了。我前面把这个值做成12,是综合了下游ClickHouse的并发写入能力、消费者线程数综合权衡的结果,并不是拍脑袋定的。分区数太少,消费者并行度上不去;分区数太多,Kafka的leader切换和元数据管理成本又上涨。在大数据场景下,一个比较稳妥的参考公式是分区数不小于消费者总数,同时也别超过broker数量的倍数太多。3个broker,12个分区,每节点4个分区,负载很均衡。
default.replication.factor=2 是副本数。日志数据量太大,3副本的存储成本有点扛不住,2副本可以容忍单节点宕机不丢数据,成本比3副本低三分之一。如果对数据安全性要求极高,可以上3副本,但成本要自己衡量。
3.2 Kafka消息延迟高的根因排查
日志消息延迟高是一个高频问题,热词里也有人专门关注"kafka消息延迟高"这个痛点。从我的经验来看,延迟高通常有几个根因:
第一是分区数不足。消费者并行处理能力受限于分区数,一个分区只能被一个消费者线程消费。假设你有3个消费者实例在消费,但topic只有3个分区,那最多3个线程并行;如果流量是100M/s,但消费能力只有30M/s,消息肯定越积越多。第二就是消费者拉取频率设置不合理,fetch.min.bytes 默认配置过大会让单次拉取的等待时间边长。第三是Kafka服务端的num.network.threads 和 num.io.threads 默认参数在高并发场景下不够用。
有一次线上通知延迟从几十秒涨到五分钟,排查过程就是按照链路一层层定位:先在消费者端打印消息接收时间戳,确认消息确实到达但消费速度跟不上;然后看Kafka的JMX指标,发现某个broker的BytesInPerSec远高于其他broker,典型的leader不均衡;最后是调整分区策略并重新分配leader。调优之后延迟恢复了秒级。
3.3 Kafka可视化工具和集群巡检
日常运维没有一个顺手的管理端工具会很吃力。我用过几款,简单聊下感受。
Kafka官方的命令行工具功能最全,但敲命令效率低;kafka-ui是开源的Web管理界面,能看topic列表、consumer group的lag,还能直接查看消息内容,日常够用;kafdrop和kafka-ui类似,界面更轻量;如果是用云厂商的托管Kafka,一般自带控制台,功能会更完善。实际运维我通常是命令行加kafka-ui配合使用——命令行做复杂操作加临时排查,kafka-ui做日常的lag监控和topic管理。
我踩过的坑是用kafka-ui直接修改了topic的分区数,导致partitions不均衡,消费延迟暴涨。后来这类操作一律走命令行和脚本,GUI工具只做查看不做变更,减少了人为事故。
4. ClickHouse表设计、建表实操与查询优化
ClickHouse是整个链路里最需要认真设计的环节。数据落的好不好,直接影响查询效率和存储成本。这块我把建表过程和踩过的坑都讲透。
4.1 核心建表语句:本地表与分布式表
ClickHouse的分布式表是逻辑表,实际数据存在本地表上。在日志场景里,我建议都采用"分布式表加本地表"的两层结构,这样既支持集群水平扩展,又能在单机层面控制数据分布。
sql复制CREATE TABLE taoke_app_log_local (
app_name String,
log_type String,
user_id String,
event_id String,
device_id String,
channel String,
referer String,
ip String,
message String,
log_time DateTime,
ingest_time DateTime DEFAULT now(),
request_id String,
order_id String,
amount Decimal(18,2)
)
ENGINE = ReplacingMergeTree()
PARTITION BY toYYYYMMDD(log_time)
ORDER BY (log_time, user_id, event_id)
TTL log_time + INTERVAL 90 DAY
SETTINGS index_granularity = 8192;
ORDER BY 是关键中的关键。在ClickHouse里它不只是排序,还承担着主索引的角色。我把user_id放在第二列,就是为了后续用户维度查询能快速定位。如果某个查询场景经常按渠道和时间过滤,把channel放到order by前面会更合适。这里要提醒的是,order by字段顺序对查询性能影响极大,一定要根据最频繁的查询模式来设计。
PARTITION BY toYYYYMMDD(log_time) 按天分区,方便数据管理。但我们实际踩过一个坑:如果某天的日志量过大,单分区数据量太大会导致查询变慢和merge压力变大。后来我们优化成按小时分区,代价是分区数量变多,merge的开销也随之上涨。这里没有绝对的对错,只有权衡。
ReplacingMergeTree 这个引擎在日志去重场景下很有用。因为Filebeat加Kafka的链路是at-least-once语义,消息可能重复。ReplacingMergeTree在后台异步去重,能在不影响写入性能的前提下处理重复数据。
4.2 分布式表与写入链路的实现
sql复制CREATE TABLE taoke_app_log_all (
app_name String,
log_type String,
user_id String,
event_id String,
device_id String,
channel String,
referer String,
ip String,
message String,
log_time DateTime,
ingest_time DateTime,
request_id String,
order_id String,
amount Decimal(18,2)
)
ENGINE = Distributed(cluster_name, default, taoke_app_log_local, rand());
分布式引擎的最后一个参数rand()是数据分布策略。用随机分布的好处是写并发均衡,但坏处是用户的日志会散落在不同分片。如果查询经常带user_id条件,可以考虑把这个参数改成cityHash64(user_id),让同一个用户的数据落在同一分片上,查询时能减少跨分片的数据拉取。我实际业务里用户查询特别频繁,所以用的是cityHash64,把同用户的数据尽量聚到同一个分片。
数据从Kafka到ClickHouse,我们最初用的是自研的消费程序,后来为了减少运维组件,尝试过ClickHouse自带的Kafka引擎表。但实际用下来,Kafka引擎表的功能弱一些——它只能做简单的物化视图落地,复杂的数据清洗、分流、多topic合并这些操作都做不了。最终方案还是保留了一个轻量级的消费程序,从Kafka拉数据、做解析和字段映射、批量写入ClickHouse。这个消费程序用Go写的,才几百行代码,维护成本很低,但功能灵活性至少翻了几倍。
批量写入参数也值得记一下:ClickHouse HTTP接口每次写入1万到5万行,缓冲时间2到3秒或者积压到5万行触发一次flush。这个参数组合在压测中效果最好,有效降低了ClickHouse的merge压力。
4.3 查询分析实操与SQL示例
试一个运营经常用到的查询——统计最近一小时各渠道的下单转化率。
sql复制SELECT
channel,
countIf(log_type = 'click') AS click_cnt,
countIf(log_type = 'order') AS order_cnt,
round(order_cnt / click_cnt, 4) AS conversion_rate
FROM taoke_app_log_all
WHERE log_time >= now() - INTERVAL 1 HOUR
GROUP BY channel
ORDER BY conversion_rate DESC;
列式引擎跑这种聚合查询非常快,几亿条数据扫描通常在一秒内完成。再比如客服要查某个用户在某个时间段的全部操作记录:
sql复制SELECT log_time, log_type, event_id, order_id, amount, message
FROM taoke_app_log_all
WHERE user_id = 'u_10012345'
AND log_time >= '2024-06-01 00:00:00'
AND log_time < '2024-06-08 00:00:00'
ORDER BY log_time
LIMIT 200;
因为ORDER BY里有user_id,这个查询会直接定位到user_id对应的索引块,加上时间范围过滤,基本是毫秒级返回。
ClickHouse查询性能问题,90%以上出在表结构设计上。我见过最典型的问题是有人把message这类长文本字段放在了ORDER BY前面,导致索引体积巨大、查询性能全面失控。长文本字段尽量别进ORDER BY,要用就放到字段列表末尾。
4.4 ClickHouse单机写入性能与硬件的关系
ClickHouse性能上限在很大程度取决于硬件。单机写入吞吐达到5万行每秒以上,磁盘建议上SSD或NVMe。我用普通SATA SSD,单机写入大概在8万行每秒;换NVMe之后能到15万到20万行每秒。内存方面,32G起步,如果并发查询多建议64G。ClickHouse的memory limit设置要根据机器实际内存调整,别用默认值太保守了,否则大查询直接OOM。
冷热数据分层的策略我也简单说下。日志数据热度集中在最近7天,查询频率最高的也是这一波。我们的做法是ClickHouse只保留最近90天数据,之前的归档到冷存储或者直接删掉。PB级日志不可能全部实时在线查询,该舍弃的就舍得舍弃。
5. 常见问题与排查技巧实录
在整套链路的运行过程中,我沉淀了几个高频问题的排查方法,直接列个速查表方便对照:
| 问题现象 | 可能原因 | 排查思路 | 解决方案 |
|---|---|---|---|
| Kafka消息堆积越来越大 | 消费者处理能力不足 | 查看consumer lag指标 | 增加消费者实例,调整fetch参数 |
| 写入ClickHouse报错too many parts | 单分区写入过快,merge跟不上 | 查看system.parts表统计 | 调节批量写入参数,减少insert频率 |
| ClickHouse查询突然变慢 | 分区内数据量过大或未按索引查 | 查看EXPLAIN语句执行计划 | 优化过滤条件,调整ORDER BY字段顺序 |
| 日志数据有重复 | Filebeat或Kafka重启重发 | 用event_id字段查重 | 使用ReplacingMergeTree去重 |
| Filebeat采集延迟高 | 日志文件滚动不及时或scan频率低 | 查看filebeat日志和偏移量 | 调整scan_frequency和close_inactive |
| 消费者频繁rebalance | 会话超时或者心跳频率设置不当 | 查看group的rebalance记录 | 调整session.timeout.ms和heartbeat.interval.ms |
除了表格里的这些,还有两个比较隐蔽的问题值得单独讲。
第一个是ClickHouse批量写入时的Too many parts报错。第一次遇到这个问题我排查了很久,后来发现就是写入频率太高,ClickHouse后台的merge线程来不及合并数据块,parts数量超过阈值就直接报错。解决方案是降低批次频率、增大单批数据量,另外调大parts_to_throw_insert的阈值只是治标不治本。
第二个是数据写入量的估算。在设计分区和集群规模时,一定要从实际数据量出发。我见到过有人用几条测试数据在单机跑通后,直接上手几十亿条线上数据,结果ClickHouse直接内存溢出,整个集群瘫痪。预研阶段至少要压到真实量级的十分之一再做容量评估,不然上线必踩坑。
还有一个小建议,日志链路一定要做数据质量监控。我在消费程序里加了几个简单的计数器,每五分钟上报一次Kafka消费的消息总数、ClickHouse写入的行数、解析失败的条数。哪边的数量对不上,直接能定位到是哪一环丢了数据。这种基础的可观测能力,在PB级场景下比任何复杂的工具都管用。
6. 关于这套方案的个人经验总结
做了这个大半年,我对这套Filebeat、Kafka、ClickHouse链路的整体评价是:简洁、稳、扩展性好。它不像那种一上来就上一堆重组件的大数据平台,反而在这种中等团队、海量日志、高实时性的场景里非常合适。
最后再分享一个我在扩容时的体会。这套链路每层都是独立扩展的——日志量涨了,加Filebeat实例;Kafka吞吐不够,加broker节点;ClickHouse查询慢,加数据分片。各层之间的解耦做得非常干净,扩容时不需要停机,一次加节点,Kafka能自动做partition的均衡,ClickHouse加节点后数据会逐步迁移。只要你把topic的partition数提前规划好,后面扩容就是加机器的事。
还有一个小坑要提醒一下:扩容ClickHouse节点之后,旧的分布式表可能还在往老节点写数据,新节点要等数据重新均衡才能发挥作用。我当时就是扩容后没有重新创建分布式表,导致新节点一直没有数据进来,白花了几天时间排查。所以扩容完记得检查分布式表的weight配置。
这套方案不是银弹,但对我们这种业务的日志分析需求来说,它是性价比最高、稳定性最可靠的组合。如果你的业务也面临类似的海量日志实时检索难题,这套架构值得一试。上面所有配置都是我实际线上在用的,照着搭就可以跑通,再根据自己环境微调参数就能用得很舒服。
