前阵子有个做实时数仓的哥们儿突然微信找我,说他们那套Kafka集群在大促期间end-to-end延迟经常飙到几秒,消费Lag曲线跟过山车似的,问我到底要不要换消息中间件。这个话题我这两年没少被问到。在很多团队里,Apache Kafka几乎就是“消息中间件”的代名词,但云原生架构铺开以后,Apache Pulsar这种存储计算分离的新设计越来越频繁地出现在技术文章、架构评审甚至面试题里。如果你正好在做技术选型,或者想把这两款中间件的原理和实战一次搞清楚,这篇内容会很对胃口。
我会从架构差异讲起,把Kafka和Pulsar在吞吐、延迟、扩容模型上的本质区别掰开揉碎,再给出一套可以直接抄作业的部署、压测、调优参数,最后附上我踩过的一些坑和面试里的高频考点。不管你是刚接触消息中间件的开发,还是已经在维护大规模集群的工程师,这篇文章都能让你少走不少弯路。
1. 为什么这两年大家开始重新审视消息中间件选型
1.1 Kafka的“光环”与扩容之痛
Kafka能在过去十几年里成为事实标准,靠的是“日志即数据”的极简设计。每个Topic拆成多个分区,分区内部追加写日志,消费者按offset顺序拉取,这套模型让它在日志收集、用户行为埋点、异步解耦这些场景里表现得非常出色。LinkedIn当初设计Kafka时就没打算把它做成通用的消息队列,而是当作一个高性能的提交日志存储来用,后来才慢慢演化出消息队列、事件流平台这些定位。
但真正维护过Kafka的人都知道它的尴尬:Broker节点既要管分区副本的同步和选举,又要负责磁盘读写和网络传输,存储和计算完全耦合在一起。集群想要扩存储,必须整节点地加机器,数据要重新平衡;某个分区因为流量倾斜成了热点,负载也只能压在那一台Broker上,其他机器帮不上忙。分区数多了以后,副本同步、文件句柄、选举开销全部上来,集群稳定性对运维的要求会越来越高。
我印象很深的一次场景是某业务Topic数量涨到500多个、分区数3000个左右时,集群里单台Broker上的分区副本数量差异已经很难控制。每次有节点滚动重启,UnderReplicatedPartitions这个指标都要抖动十几分钟,ISR频繁收缩。当时所有人都知道这种架构迟早会碰到上限,但考虑到业务迁移动成本太大,只能硬着头皮加机器。
1.2 Pulsar重新定义了“存储”
Pulsar的思路其实更接近云原生里面常说的“无状态”:计算和存储分层。Broker节点只管消息的路由、订阅管理和IO请求调度,本身不保存数据;实际的数据持久化全部交给底层的Apache BookKeeper集群处理。BookKeeper里的Bookie节点可以单独扩展存储容量,Broker可以单独扩展计算能力,两者互不拖累。
这种分层带来的第一个直接好处是扩容变“优雅”了。传统Kafka加节点后要等数据再平衡,期间业务还可能受影响;Pulsar这边加一个Broker,只要把负载均衡策略打开,Topic的分区会自动往新节点上迁移,数据层完全不用动。反过来,如果只是存储快满了,加几台Bookie节点就能把容量顶上。
另外一个杀手级特性是分层存储。Kafka的数据不管多冷都得放在Broker本地磁盘,除非你自己写脚本做归档;Pulsar原生支持把老数据、冷数据卸载到对象存储或者HDFS上,Topic可以保留完整的数据历史,但热数据只占一小块高性能盘,长尾存储成本大幅下降。
我之前在一个实时数仓项目里做过粗略估算,同样的数据保留周期,Kafka方案如果全部用SSD可能是Pulsar本地Bookie加对象存储冷备方案的3到4倍成本。当然这不是说Kafka不好,而是不同架构在不同业务模型下注重的指标本来就不一样。
2. 架构层面的核心分野:存储与计算的关系
2.1 Kafka的分区日志模型
为了讲清楚Kafka的性能边界,还是得把它的存储模型拉出来仔细看。Kafka里一个Topic被切分成多个Partition,每个Partition就是一份有序的日志Segment文件序列。生产者把消息追加到Leader分区日志的末尾,Follower副本从Leader拉取数据做同步,消费者按照提交过的Offset顺序拉取消息。
这套模型之所以快,主要靠三个杀手锏:顺序写、Page Cache、零拷贝。写消息时多个生产者批次数据连续追加,磁盘顺序写的速度在现代NVMe盘上可以达到每秒GB级别;读数据时操作系统Page Cache会把热数据留在内存里,消费者拉取请求很大概率直接命中缓存,根本不走磁盘;零拷贝则能把内核态到用户态的数据拷贝省掉,把网络吞吐拉到接近网卡极限。
不过这个设计的短板也很清楚:一个分区在任意时刻只能被一个消费组里的一个消费者读取。想提高消费并发,只能把分区数加多,但分区数不是无限涨的——每个分区都要有Leader和多个Follower,都占用Broker上的文件和线程资源,副本同步也要吃网络带宽。分区数过多时,Broker的性能反而会明显下降。
2.2 Pulsar的Broker + BookKeeper分层模型
Pulsar把消息的“路由”和“存储”拆成了两个完全独立的层。客户端通过Service Discovery找到Topic所在的Broker,Broker负责维护元数据、订阅游标、消息的读写调度,但实际消息数据被切成多个Ledger分片,写入BookKeeper中的一组Bookie节点。
BookKeeper的存储模型要稍微花点时间消化。每个Ledger被拆成多个Segment,每个Segment会在多个Bookie上写多份副本,这套机制有点类似“分布式日志数组”。写入时Bookie先落Journal日志保证故障恢复时不丢数据,再把数据刷到EntryLog文件里,同时配合读缓存提高热数据的读取效率。因为Broker完全不碰磁盘,它就可以做到无状态水平扩展,Topic的分区可以动态分布在任意多个Broker上。
更关键的是,Pulsar的订阅模型跟Kafka的消费组模型不是一个路子。Kafka里一个分区只能被消费组里的一个消费者独占;Pulsar除了Exclusive独占模式,还有Shared和Key_Shared等模式,同一分区内的消息可以让多个消费者并行处理,这在实时流计算场景里非常有用,等于把单分区内的消费并发度也释放出来了。
2.3 从吞吐和延迟的角度看差异
从Benchmark的角度说,Kafka在小规模日志型场景下确实非常能打,单Broker顺序写吞吐能到几十万条每秒,端到端延迟用批量发送控制好也能压到几十毫秒甚至更低。它的问题更多出现在“规模变大”之后:分区数量一多,副本同步、Leader切换这些分布式协调操作都会变成瓶颈,扩容时的数据重平衡也很痛。
Pulsar因为存储计算分离,理论上Broker不存数据,所以Topic数量可以做得非常大,社区里甚至有管理十万级Topic的案例。每增加一个Topic,Pulsar只需要在元数据层登记一下,数据分片会均匀地散到BookKeeper节点上,这跟Kafka单个Broker要同时顶着多个分区副本的负载是完全不同的感受。很多人容易忽略的是,当分区规模远超节点数时,Pulsar的“隔离性”优势会被放大得非常明显。
不过Pulsar也不是银弹。多了一层BookKeeper之后,链路变长,端到端延迟的抖动因素也变多了。BookKeeper的Journal盘如果不够快,写入延迟会直接传导给生产端;Broker和Bookie之间的网络如果跨机房,延迟也会显著增加。所谓“Pulsar比Kafka快”这种说法并不严谨,它更准确的定位是“在同等规模下更容易线性扩展,在更大规模下单分区性能上限更高”。
3. 从零搭建集群并压测出真实性能
3.1 硬件与部署形态的选择
先说硬件。Kafka和Pulsar对磁盘的偏好不太一样。Kafka因为数据全部在本地,建议直接把日志目录放到NVMe SSD上,至少也要是高性能SATA SSD,普通HDD只能用来做容量型日志场景,延迟指标会比较难看。Pulsar这边更讲究分工——Bookie节点的Journal盘建议用NVMe SSD承担写入,EntryLog数据盘可以用大容量SSD或HDD,容量和性能分开算,别一根筋地把所有数据都塞进一块盘。
部署形态上,Kafka 3.x之后已经可以用KRaft模式替换掉ZooKeeper,集群成员管理、元数据一致性都交给Controller节点处理,省掉一套ZK运维负担。Pulsar本身是围绕ZooKeeper、BookKeeper、Broker三个核心组件运转的,ZooKeeper负责元数据和协调,BookKeeper负责存储,Broker负责计算接入,生产环境通常还会加一层Proxy用于客户端接入和安全管控。
3.2 Kafka集群部署核心配置
我这边以Kafka 3.6 KRaft模式三节点集群为例,直接挑核心配置讲。控制器节点和Broker节点可以合并角色,只是生产环境我更推荐把Controller独立出来,让元数据层不要被数据流量干扰。
# config/server.properties 核心项 process.roles=broker,controller node.id=1 controller.quorum.voters=1@192.168.1.11:9093,2@192.168.1.12:9093,3@192.168.1.13:9093 listeners=PLAINTEXT://192.168.1.11:9092,CONTROLLER://192.168.1.11:9093 advertised.listeners=PLAINTEXT://192.168.1.11:9092 log.dirs=/data/kafka-logs num.partitions=12 default.replication.factor=3 min.insync.replicas=2注意几个容易踩的坑。advertised.listeners一定要填客户端能访问到的地址,很多内网IP问题都出在这;log.dirs可以配多个目录,但最好让同一副本的不同目录分散在不同物理磁盘上,别全放一起;min.insync.replicas是“至少几个副本确认才算写入成功”的兜底,设成2就是为了配合生产端acks=all实现不丢消息,如果只写1个副本,Leader宕机就可能丢数据。
分区数这块,建议按照消费并发和单分区带宽预期来定。我自己常用的经验公式是:单分区吞吐按5~20MB/s预估,目标Topic峰值流量除以这个区间,再乘以副本因子和未来增长系数。比如一个日志主题峰值300MB/s,按15MB/s单分区算,20个分区起步,副本3份,再留30%余量,定到32个分区比较稳。
3.3 Pulsar集群部署核心配置
Pulsar用Helm在Kubernetes上部署是最常见的。组件多,但核心就几个:ZooKeeper、Bookie、Broker,再加上可选的Proxy。我建议先在单机装上standalone熟悉一遍操作。
# 单机体验 bin/pulsar standalone # 生产环境Bookie关键配置 journalDirectory=/data/pulsar/journal ledgerDirectories=/data/pulsar/ledgers zkServers=zk1:2181,zk2:2181,zk3:2181Broker上最核心的容错参数是这三个:managedLedgerDefaultEnsembleSize、managedLedgerDefaultWriteQuorum、managedLedgerDefaultAckQuorum。EnsembleSize决定消息分片写到多少个Bookie上,WriteQuorum决定写几个副本,AckQuorum决定收到几个确认才返回成功。常规三副本配置就是EnsembleSize=3, WriteQuorum=3, AckQuorum=2,既能容忍一个Bookie故障,又不会因为必须写满三个副本而拖慢延迟。
Bookie磁盘分配上,Journal目录和Ledger目录必须分开,Journal盘优先用延迟极低的NVMe盘,Ledger目录可以用大容量SSD甚至HDD。这样即使Ledger盘写入慢,Journal也能快速确认写入,再异步刷到Ledger盘,生产延迟的抖动会被控制在小范围内。
3.4 压测工具与结果解读
Kafka自带两把压测利器,直接就能用。
# 生产者压测:500万条,每条1KB,acks=all bin/kafka-producer-perf-test.sh \ --topic perf-topic \ --num-records 5000000 \ --record-size 1024 \ --throughput -1 \ --producer-props bootstrap.servers=192.168.1.11:9092 \ acks=all linger.ms=5 batch.size=65536 # 消费者压测 bin/kafka-consumer-perf-test.sh \ --bootstrap-server 192.168.1.11:9092 \ --topic perf-topic \ --messages 10000000 \ --threads 8Pulsar自带的pulsar-perf命令更清晰。
bin/pulsar-perf produce -s 1024 -n 1000000 -t 10 -r 1000 perf-topic bin/pulsar-perf consume -s 100000 -n 1000000 perf-topic压测时别只看吞吐量,延迟分位数比平均值重要得多。P99如果稳稳压在几十毫秒,同时P999没有明显掉队,才是健康状态。还有一点是务必按维度拆压测:先单分区、单Topic测出基准值,再整体压测评估集群容量,不然瓶颈到底出在磁盘、网络还是锁竞争都定位不清楚。
4. 低延迟实时处理的关键调优
4.1 端到端延迟到底由什么组成
所谓“低延迟”不能只看生产端写入快慢,一条消息从Producer发送到Consumer消费完成,经过的环节是:客户端网络发送、Broker日志写入或BookKeeper持久化、Consumer网络拉取、反序列化和业务处理。每个环节都可能变成瓶颈。
Kafka端到端延迟的常见模型是:生产者批量攒数据(linger.ms)加上Broker刷盘时序加上消费者poll间隔。Pulsar端到端延迟还要多算上一跳BookKeeper写入确认。很多人在Kafka里调了好几天,最后发现瓶颈在消费端的反序列化GC停顿上,这种情况换中间件根本解决不了。
4.2 Kafka低延迟优化清单
生产端追求低延迟,核心思路是“别攒太久,但又要利用批处理降低网络往返”。linger.ms=0时每条数据都立即发送,延迟最低,但吞吐会打折扣;业务上如果允许微秒级延迟,建议把linger.ms设成1~5,批次大小batch.size调到64KB左右,既能合并小消息,又不至于为了凑批牺牲延迟。compression.type开lz4或zstd,网络传输量下降后延迟也会跟着下降,代价是消耗一点CPU。
Broker端需要关注线程模型。num.network.threads一般设为CPU核数的两倍,负责处理网络请求;num.io.threads设为CPU核数,负责执行实际磁盘读写和副本拉取。如果CPU核数多但磁盘是机械盘,num.io.threads拉太高反而会让线程排队加剧,需要结合监控实测来定。
Consumer端的低延迟容易被人忽略。fetch.min.bytes=1表示Broker一有数据就返回,fetch.max.wait.ms控制在500以内,这样消费者拉取请求不会干等数据攒批。还有一个更隐蔽的点:max.poll.records别设太大,一条一条或少量处理完再去poll,既能降低单次心跳超时风险,也能减少长轮询对延迟的拖累。
4.3 Pulsar低延迟优化清单
Pulsar调优要同时看Broker和Bookie两层。
Broker端要开消息批次功能。batchingEnabled=true、batchingMaxMessages=1000、batchingMaxPublishDelay=1,这个“攒1毫秒再发”的机制跟Kafka的linger.ms思路一样,能在延迟几乎无损的情况下把吞吐拉高一个量级。Consumer端把receiverQueueSize从默认的1000适当提高,可以抵消网络往返的开销。
Bookie端最容易出问题的点是Journal刷盘策略。生产环境如果对延迟极度敏感,可以把journalSyncData关掉,让Journal写入异步刷盘,换取更低的写确认延迟,但代价是极端宕机场景可能丢最后一批数据。这个开关一定要跟业务方确认“最多一次”还是“至少一次”语义,别自作主张改。
另外Broker的堆内存设置也必须留足,managedLedgerCacheSizeMB默认是堆内存的40%,如果堆开太大但缓存占比不够,读热点就会频繁穿透到Bookie,延迟长尾就来了。
4.4 消息延迟突然变高的排查实战
有一次线上Kafka集群突然出现部分Topic延迟飙高,消费Lag涨到几百万,但CPU、磁盘、网络看起来都还算正常。我先看消费端,发现有个消费者实例的GC从正常的200ms变成秒级停顿,Full GC频繁到几分钟一次。原因很快定位:消费端把每条消息里的JSON体展开成一个超大对象图,堆内存被快速打满。
换掉序列化方案、调整堆参数之后GC恢复正常,但这只解决了“消费端”的问题。接着排查生产端,发现有个业务方为了拿更低延迟把linger.ms设成了0,但消息体平均有几十KB,小包高频发送让Broker端的网络线程几乎满载。把生产端的批处理和压缩打开后,Broker CPU立刻降下来,端到端延迟回到几十毫秒。
排查这类问题一定要按“生产端 -> Broker -> 消费端”的顺序逐层看。先看生产端的发送TPS和错误率,再看Broker的BytesInPerSec、UnderReplicatedPartitions、请求队列堆积长度,最后用kafka-consumer-groups.sh --describe看每个分区的Lag分布。哪一层数据异常,瓶颈基本就锁定在哪一层。
5. 面试与选型的高频考点速查
5.1 两分钟讲清楚架构差异
如果你面试时被问到“Kafka和Pulsar有什么区别”,别一上来丢一堆术语,先用一句话提纲挈领:Kafka把存储和计算绑定在同一个Broker上,靠分区多副本提升吞吐;Pulsar把计算(Broker)和存储(BookKeeper)分层,Broker无状态,存储节点可独立扩展。然后展开三点:部署形态、扩容模型、消费模型。
扩容模型是最容易让面试官眼睛一亮的点。Kafka扩容要数据重平衡,分区会重新分配,期间IO和网络压力会上升;Pulsar加Broker只需要做负载均衡迁移Topic归属,数据层不动。消费模型上,Kafka一个分区只能被同一消费组里的一个消费者消费,Pulsar支持Shared和Key_Shared模式,单分区多消费者并行。
5.2 面试中常见的Kafka/Pulsar问题
我整理了几道几乎每次都会被问到的题,附上我自己的答题思路。
- Kafka为什么能支撑高吞吐?答:顺序写磁盘、Page Cache命中、零拷贝、批量发送和压缩,四个点一起说才完整。
- Kafka的ISR机制是什么?答:ISR是跟Leader保持同步的副本集合,Producer的acks和min.insync.replicas配合决定了消息持久化强度,ISR收缩后要能及时察觉并修复。
- 消息堆积怎么处理?答:先分场景。如果是消费者处理慢,扩消费者实例或者调大分区并行度;如果是分区热点导致某个分区堆积,要做消息key的散列策略或者拆分Topic;如果是下游依赖响应慢,考虑异步批量写下游。
- 如何保证消息不丢?答:生产端acks=all加重试,Broker端min.insync.replicas=2,消费端关闭自动提交或手动确认,每一步都要明确至少一次还是精确一次。
- Pulsar相比Kafka订阅模型有什么优势?答:共享订阅支持单分区多消费者并行,Key_Shared按key路由保证同一key有序,这在流式计算里很实用。
- Pulsar的BookKeeper写入模型怎么理解?答:类比成分布式日志分段存储,消息拆成Ledger和Segment,每个Segment写多个Bookie,采用Quorum机制确认,Journal先落盘保证恢复能力。
5.3 最终选型速查表
| 对比维度 | Kafka | Apache Pulsar |
|---|---|---|
| 架构模型 | 存储计算耦合 | 存储计算分离 |
| 集群扩容 | 需要数据重平衡 | Broker和Bookie独立扩容 |
| 单分区消费并发 | 一个分区只能被消费组内一个消费者消费 | Shared/Key_Shared支持多消费者并行 |
| Topic数量规模 | 大Topic数场景运维压力大 | 支持超大Topic数 |
| 冷数据存储 | 依赖本地磁盘,需自行归档 | 原生分层存储到对象存储 |
| 生态成熟度 | 生态极丰富,周边工具多 | 生态在快速成长,部分功能需要自己封装 |
| 运维复杂度 | 单层组件,概念少 | ZK+Bookie+Broker+Proxy,概念多 |
| 典型场景 | 日志管道、大数据生态、实时数仓 | 多租户平台、超大规模事件流、弹性云原生 |
严格说这两者不是“谁替换谁”的关系。Kafka在生态对接上的优势太强了,Flink、Spark、ClickHouse周边组件几乎默认支持Kafka协议;Pulsar则更适合那种对弹性、多租户、存储成本和分区规模有极致要求的平台型系统。
6. 聊聊从Kafka迁移到Pulsar的真实体会
6.1 什么场景值得换
我见过不少团队听说Pulsar能解决扩容痛点就想立刻搬,但我得泼盆冷水:如果你们当前的Kafka集群规模不大,Topic数量几百个,分区几千个,日常没有明显的扩容痛苦,迁移的收益其实不大。迁移本身要重写客户端配置、重验证消费语义、重做监控告警,光这些人力成本就够维护好几年Kafka了。
反过来,如果业务有这几个特征,Pulsar的吸引力会非常明显:一是Topic数量增长极快,半年内可能翻好几倍;二是需要多团队共用一套消息平台,要严格的多租户隔离;三是有大量历史数据要保留但不想承担全量SSD成本;四是集群需要频繁弹性伸缩,比如每天不同时段的流量差异非常明显。这些场景下,Kafka的运维成本会随规模非线性上涨,Pulsar的架构优势才能真正兑现。
6.2 迁移前你必须想清楚的几件事
跨中间件迁移最大的坑是消费语义对不上。Kafka用Consumer Group加Offset管理消费进度,Pulsar用Subscription和Cursor管理,虽然都能实现“至少一次”,但“精确一次”的落点完全不同。Flink任务切换消息源时,Checkpoint对齐机制要重新验证,否则可能出现重复或丢失。
迁移顺序我建议走“双跑”策略:新集群和老集群并行运行一段时间,业务写入两边,消费端先切一部分流量到新集群验证,确认Lag、延迟、稳定性都符合预期后,再逐步放量切完。不要想着周末一个窗口把所有消息一把梭,消息中间件的迁移是最不值得赌的。
最后再分享一个小技巧:Kafka客户端自带一个kafka-reassign-partitions.sh,可以在不停业务的情况下慢慢把分区挪到新扩的节点上,但要注意把reassignment.throttle限流参数调低一点,不然网络带宽很容易被打满,业务流量跟着遭殃。Pulsar这边做Broker负载均衡时也有类似的流量限制开关,生产环境一定不要贪快。
选型这件事,到最后拼的不是哪个中间件更“高级”,而是哪个中间件在你真实的业务模型里更省心。我自己的体会是:先花一天时间把吞吐量、Topic规模、消息保留策略、消费并发模型这些数字摸清楚,再决定要不要换,绝对比盲目追新兴架构划算得多。