IT 论文精读 · PAPER 58
Kreps, Narkhede & Rao · LinkedIn · NetDB 2011
2011 年,LinkedIn 的三位工程师发了一篇只有七页的短论文,介绍他们自造的 Kafka。它是公司内部的一条数据主干道:所有服务器产生的动静(谁点了什么、谁搜了什么、哪台机器 CPU 高了)都往这条道上扔,需要的人再从道上自取。今天从电商的订单状态到银行的实时风控,背后大多跑着 Kafka 或它的模仿者。
稍大的网站每天产生的行为流水——每一次点击、刷新、划过——比订单账户这类「正经数据」多好几个数量级。当年处理它只有两条路。
一条是传统消息队列,像个讲究的邮局:每封信挂号、签收、登记谁取走了,取走即销毁。稳妥,可光记这些账就贵得扛不住量。另一条是日志搬运工:每台机器把日志攒着,定时打包搬进大仓库集中算——量扛得住,可搬一趟要几小时,等算出「这人可能想看什么」,人早走了。
要快的没有量,有量的不够快。
Kafka 反过来:服务器不再替任何人记账。
它把消息写成一本一直往后翻的流水账——来一条就在末尾添一行,从不修改、从不插页。谁想读,自己记住读到第几行,下次接着往下读。读过的内容不会消失,本子按时间到期(比如七天)才整段扔掉。
于是「一条消息只能被取走一次」这条天经地义的规矩没了:搜索的读它、推荐的读它、报表的半夜读它,各拿各的书签、互不打扰。程序出了 bug?书签往回拨几页重读一遍就是。
一是只往末尾添行。硬盘最怕东一榔头西一棒子地找地方,最擅长一路往下写。
二是服务器不记账。「谁读到哪儿」原本是最烦的一本账,每条消息都要记状态、还要建索引翻查。丢给读者自己保管后,服务器只剩两件事:往后添、按行号取。
三是取的时候不绕路。读者说「从第几行起给我一段」,服务器就把那段字节直接推进网线,不先搬进自己内存倒一手。
四是一本拆成多本。同一话题的账本拆成几本、分放不同机器,就能同时写、同时被读。
诚实说一句代价:论文里的 Kafka 每条消息只存一份,那台机器硬盘彻底坏了,还没被读走的消息就永远没了(多存几份是后来才补的)。
「日志」从没人管的副产品变成整家公司数据的主干道:一份流水摆在那儿,谁想用就接上去读,新接一个系统不必再惊动上游。今天但凡讲「实时数据」「事件驱动」的系统,大多照这个样子搭。
把消息中转站从「替人记账的邮局」改成「一本只往后添、按期作废的流水账」:服务器不管谁读到哪儿、书签由读者自己拿——于是又快、又能多方共读、还能倒带重放。
想看它的架构图、offset 与零拷贝机制、以及实测数字? → 切到精读版
Kafka 把消息系统重做成一个只追加的分布式日志(append-only log):消息按主题(topic)切成若干分区(partition),每个分区就是一个不断往末尾追加的日志文件;服务端(broker)不记录谁消费到哪儿、也不因消息被读走而删它(只按时间到期整段丢弃),进度由消费者自己保管。这套「无状态 broker + 顺序读写 + 拉取」让普通机器达到每秒几十万条写入,也让同一份数据被多方各按节奏消费、还能倒带重放。
作者是 LinkedIn 的 Jay Kreps、Neha Narkhede、Jun Rao,论文发在 NetDB'11(2011 年 6 月与 SIGMOD 同期在雅典举办的小型工作坊),正文只有七页。它明说自己是把两条既有路线的好处合起来:企业消息系统(IBM WebSphere MQ、JMS 及其实现 ActiveMQ、RabbitMQ)与日志聚合器(Facebook 的 Scribe、Cloudera 的 Flume、Yahoo 的 Data Highway)。Kafka 2011 年开源、次年成为 Apache 顶级项目;三位作者后来创立了 Confluent。
互联网公司每天产生海量日志数据:登录、浏览、点击、搜索,加上调用延迟、CPU 等运维指标。论文点出的关键变化是——它过去只是事后分析的素材,如今直接进了线上功能:搜索相关性、推荐、广告定向、反垃圾、信息流。而它比订单账户这类「正经数据」大几个数量级:光算点击率,就得为页面上每个没被点的条目也记一条。
于是需求成了「既要扛得住量、又要几秒内可用」,而两类现成系统各缺一半:企业消息系统功能过剩(逐条确认这类强保证,对「偶尔丢几条浏览事件也无妨」的日志是奢侈品),且不以吞吐为首要目标(JMS 没有批量发送接口,每条消息一个完整 TCP 往返)、分布式支持弱、还假设消息被近乎即时取走,一旦积压性能就显著劣化;日志聚合器则为离线而生,周期性倒进 HDFS 或数据仓库,天然是小时级的,且多用推(push)模型,慢的消费者会被推垮。
基本概念只有四个:一类消息的流叫一个主题;生产者向主题发布消息;消息存在一组叫 broker 的服务器上;消费者订阅主题,靠拉(pull)取走消息。为了摊开负载,一个主题被切成多个分区,每个 broker 存其中一部分分区。
每个分区在物理上是一组大小相近的段文件(segment file,例如 1GB 一个)。生产者发来消息,broker 只做一件事:追加到最后一个段文件的末尾;攒够一定条数或时间才刷盘,刷盘之后才对消费者可见。
最不寻常的一刀在这里:Kafka 的消息没有显式 ID,直接用它在日志中的位置(offset)寻址。这省掉了传统消息系统里那套「ID → 物理位置」的随机访问索引——既要维护,查起来又满盘找。代价是 offset 递增但不连续(下一条 = 当前 + 当前消息长度,本质就是字节位置)。
那怎么定位?broker 内存里只留一张很小的有序表:每个段文件首条消息的 offset。收到「从 offset X 起给我最多 N 字节」的请求,先查出 X 落在哪个段,再顺着读出去。
全篇最反直觉的一刀:broker 不记录任何消费者消费到了哪儿。好处是它少了一大堆复杂度与开销——不必为每条消息维护投递状态,也就不必为此建索引、写盘。
可这样一来 broker 不知道消息是否已被所有订阅方读走,什么时候能删?答案简单粗暴:按时间保留——超过一段时间(论文说典型是 7 天)就删,与有没有人读过无关;理由是消费者要么实时、要么按小时按天消费,而 Kafka 性能不随数据量劣化,长期保留负担得起。
这一刀还砍出一个红利:消费者可以故意把 offset 往回拨、重新消费。这违反了队列「取走即消失」的常识,却是刚需——消费端逻辑写错了,改完可以重放。论文特意点明倒带在拉模型里很容易、在推模型里很难;而「拉」本就是明确立场:消费者按自己能承受的最大速率取数据,不会被推垮。
① 不在自己进程里缓存消息,直接吃操作系统页缓存。Kafka 用 JVM 写成,却完全不做应用层消息缓存:避免双份缓冲、进程重启后缓存依然是热的、几乎没有消息对象因而垃圾回收负担极轻。更妙的是生产者与消费者都顺序访问段文件、且消费者通常只落后一点点,正撞在操作系统写透缓存与预读的枪口上。作者称生产与消费的性能都与数据量呈线性、直到许多 TB。
② 用 sendfile 省掉两次拷贝一次系统调用。把本地文件发到远端 socket 的常规做法要走四步:磁盘→页缓存、页缓存→应用缓冲区、应用缓冲区→socket 缓冲区、再上网卡,共 4 次拷贝 + 2 次系统调用。sendfile 把字节从文件通道直送 socket 通道,省掉中间那 2 次拷贝和 1 次系统调用。前提是 Kafka 不需要在中途改动消息:存的什么格式,发的就是什么格式。
③ 批量。生产者一次请求发一组消息,消费者每次拉取也是一批(典型几百 KB),RPC 的固定开销就被摊薄了。
消费者组(consumer group):组内一条消息只投给其中一个消费者(点对点),组间各自独立地拿到全量(发布订阅),组与组之间无需任何协调。两个设计决定值得记:
论文给得坦白:一般只保证「至少一次」——「恰好一次」要两阶段提交,作者认为对其场景不必要。消费者非正常崩溃时,接管其分区的人会重读一小段(已消费、但 offset 还没提交到 ZooKeeper 的那些)而产生重复,在意的应用要自己按 offset 或唯一键去重。顺序上:单分区内有序,跨分区不保证。每条消息还带 CRC,broker 遇 I/O 错误会剔掉对不上的。
生产环境。每个线上数据中心各配一个 Kafka 集群;分析数据中心另有一个集群,用内嵌消费者把线上数据拉过来,再送进 Hadoop 与数据仓库。整条管线端到端平均约 10 秒(论文说「没怎么调优」,已够用);当时的量是每天数百 GB、接近十亿条消息。
对比实验。对手是 ActiveMQ v5.4(默认存储 KahaDB)与以性能著称的 RabbitMQ v2.4;两台 Linux 机器各 8 个 2GHz 核、16GB 内存、6 块盘做 RAID 10,千兆互联,一台当 broker,都配成异步刷盘。
作者主动加了句限定:实验目的不是证明别的系统差——两个对手的功能都比 Kafka 多——而是说明专用系统能换来多大的性能空间。
Kafka 真正的贡献不是「一个更快的队列」,而是把「日志」提升成了一等的系统抽象。当消息不再随消费而消失、进度由消费者自持,「队列」与「存储」的界限就模糊了:同一份有序的事实流,实时服务紧咬队尾读,数据仓库慢半天读,新系统上线还能从头重放历史。这解开了大公司里最难缠的一团线——原本 N 个数据源各自对接 M 个下游,如今都先汇进这条主干道。论文里未竟的两件事后来成了主战场:跨 broker 复制补上持久性,流处理长成了 Kafka Streams。
① 一句话:把消息系统重做成「分区 + 只追加日志」,broker 不记消费进度、不因被读而删,消费者自持 offset。
② 痛点:行为日志比业务数据大几个数量级;企业消息系统扛不住量(逐条确认、无批量、积压即劣化),日志聚合器只能小时级离线。
③ 存储:分区 = 一串约 1GB 的段文件,只往末尾追加;没有消息 ID,用字节位置 offset 寻址,省掉随机访问索引,内存里只留「每段首条 offset」的小表。
④ 无状态 broker:进度交给消费者,删除改为按时间到期(典型 7 天);红利是可以倒带重放,而倒带在拉模型里才好做。
⑤ 快的来源:不做应用层缓存、直接吃页缓存;sendfile 零拷贝;收发都批量。协调上用消费者组,分区是并行最小单位,不设 master 而靠 ZooKeeper 自行重平衡。
⑥ 结果:单生产者 5 万条/秒(批 1)、40 万条/秒(批 50),消费 2.2 万条/秒(超对手四倍);每条消息开销 9 字节 vs ActiveMQ 的 144 字节。LinkedIn 每天数百 GB、近十亿条,端到端约 10 秒。
⑦ 保证与局限:至少一次、分区内有序;论文版本无复制、生产者不等确认、到期即删会让慢消费者永久错过、ZooKeeper 依赖与重平衡抖动。
⑧ 影响:把「日志」变成公司数据管道的中枢抽象,成了流处理与事件驱动架构的事实底座。