Kafka:面向日志处理的分布式消息系统
Kafka: A Distributed Messaging System for Log Processing
贡献:提出基于日志的分布式消息队列模型,统一了日志收集与流式数据传输的架构。,如今已成为工业界消息队列的事实标准,支撑了几乎所有实时数据处理场景。
Jay Kreps Neha Narkhede Jun Rao
LinkedIn Corp.
摘要
日志处理已成为消费级互联网公司数据管道中的核心组成部分。本文介绍了Kafka——一款我们为采集与传输高容量日志数据、实现低延迟交付而开发的分布式消息系统。该系统融合了现有日志聚合系统与消息系统的设计思想,同时适用于离线与在线消息消费场景。我们在Kafka中做出了诸多非常规但具备实用性的设计选择,以实现系统的高效性与可扩展性。实验结果表明,与两款主流消息系统相比,Kafka具备更优异的性能表现。目前Kafka已在生产环境中投入使用一段时间,每日可处理数百GB的新增数据。
通用术语
管理、性能、设计、实验验证
关键词
消息传递、分布式、日志处理、吞吐量、在线处理
1. 引言
任何具备一定规模的互联网公司都会产生海量的“日志”数据。这类数据通常包括:(1)用户活动事件,对应登录、页面浏览、点击、点赞、分享、评论与搜索查询等行为;(2)运维指标,例如服务调用栈、调用延迟、错误信息,以及每台机器的CPU、内存、网络、磁盘使用率等系统指标。长期以来,日志数据都是分析工作的组成部分,用于追踪用户参与度、系统利用率及其他指标。然而,近年来互联网应用的发展趋势,使得活动数据成为生产数据管道的一部分,直接用于网站功能特性中。这类应用场景包括:(1)搜索相关性优化;(2)基于物品热度或活动流中出现频次生成的推荐内容;(3)广告定向与效果报表;(4)防范垃圾信息、非法数据爬取等滥用行为的安全应用;(5)聚合用户状态更新或动态,供其“好友”或“人脉”查看的信息流功能。
日志数据的这种生产级、实时化应用,给数据系统带来了新的挑战,因为其数据量比“真实”业务数据高出数个数量级。例如,搜索、推荐与广告场景通常需要计算细粒度的点击率,这不仅会为每次用户点击生成日志记录,还会为页面上数十个未被点击的物品生成日志。中国移动每天会采集5~8TB的通话记录[11],Facebook每天会收集近6TB的各类用户活动事件[12]。
早期处理这类数据的系统,大多依赖从生产服务器上物理抓取日志文件来开展分析。近年来,业界已构建了多款专用的分布式日志聚合系统,包括Facebook的Scribe[6]、雅虎的数据高速公路(Data Highway)[4]以及Cloudera的Flume[3]。这些系统主要用于采集日志数据并加载到数据仓库或Hadoop[8]中,供离线消费。在领英(一家社交网站),我们发现除了传统的离线分析之外,还需要支持上述各类实时应用,且延迟不能超过数秒。
我们开发了一款名为Kafka[18]的新型日志处理消息系统,它兼具了传统日志聚合系统与消息系统的优势。一方面,Kafka具备分布式、可扩展的特性,能够提供高吞吐量;另一方面,Kafka提供了类似消息系统的API,允许应用实时消费日志事件。Kafka已开源,并在领英的生产环境中成功运行了6个月以上。它极大简化了我们的基础设施,因为我们可以用同一套软件实现各类日志数据的在线与离线消费。本文的其余部分安排如下:第2节回顾传统消息系统与日志聚合系统;第3节介绍Kafka的架构及其核心设计原则;第4节介绍领英的Kafka部署情况;第5节展示Kafka的性能测试结果;第6节讨论未来工作并总结全文。
2. 相关工作
传统企业消息系统[1][7][15][17]已存在多年,通常作为事件总线在异步数据流处理中发挥核心作用。然而,有诸多原因导致它们并不适合日志处理场景。首先,企业系统提供的功能与日志处理的需求不匹配。这类系统往往侧重于提供丰富的交付保证机制。例如,IBM Websphere MQ[7]支持事务,允许应用将消息原子性地插入多个队列。JMS[14]规范允许每条消息在消费后单独确认,且确认顺序可以与消费顺序不一致。对于日志采集而言,这类交付保证往往属于过度设计。例如,偶尔丢失少量页面浏览事件,并不会造成严重影响。这些非必要的特性,增加了系统API与底层实现的复杂度。其次,很多系统并未将吞吐量作为核心设计约束。例如,JMS没有提供API允许生产者显式地将多条消息批量打包到单个请求中。这意味着每条消息都需要一次完整的TCP/IP往返,对于我们场景下的吞吐量需求而言并不可行。第三,这类系统的分布式支持能力较弱,无法便捷地将消息分区存储到多台机器上。最后,很多消息系统默认消息会被立即消费,因此未消费消息的队列规模通常很小。当消息产生堆积时(例如数据仓库类离线应用,它们不是持续消费,而是周期性地进行批量加载),系统性能会显著下降。
过去几年间,业界推出了多款专用日志聚合系统。Facebook使用一款名为Scribe的系统:每台前端机器通过套接字将日志数据发送到一组Scribe节点,每台Scribe节点对日志条目进行聚合,再定期转储到HDFS[9]或NFS设备中。雅虎的数据高速公路项目也采用了类似的数据流:一组机器聚合来自客户端的事件,生成“分钟级”文件,再写入HDFS。Flume是Cloudera开发的一款较新的日志聚合系统,它支持可扩展的“管道”与“接收器”,让日志数据的流式传输非常灵活,同时也提供了更完善的分布式支持。然而,这类系统大多面向离线日志消费场景,且往往会向消费者暴露不必要的实现细节(例如“分钟级文件”)。此外,它们大多采用“推送”模型,由代理节点将数据转发给消费者。在领英,我们发现“拉取”模型更适配我们的应用场景:每个消费者都可以按照自身能承载的最大速率拉取消息,避免因推送速度超过处理能力而导致消息淹没。拉取模型还让消费者的回溯消费变得更容易,我们将在3.2节末尾讨论这一优势。
近期,雅虎研究院开发了一款名为HedWig[13]的新型分布式发布订阅系统。HedWig具备高可扩展性与高可用性,能够提供可靠的持久性保证,但它主要用于存储数据存储的提交日志。
3. Kafka架构与设计原则
基于现有系统的局限性,我们开发了一款基于消息模型的新型日志聚合系统Kafka。首先介绍Kafka中的基本概念:特定类型的消息流被定义为一个主题(topic);生产者可以向主题发布消息;发布的消息会存储在一组被称为**代理节点(broker)**的服务器上;消费者可以从代理节点订阅一个或多个主题,并通过拉取数据的方式消费订阅的消息。
消息传递在概念上非常简单,我们也尽量让Kafka的API保持同样的简洁性。我们不展示完整的API定义,而是通过一些示例代码说明API的使用方式。下方是生产者的示例代码。消息被定义为仅包含字节型有效载荷,用户可以选择自己偏好的序列化方式对消息进行编码。为提升效率,生产者可以在单次发布请求中发送一组消息。
生产者示例代码:
producer = new Producer(...);
message = new Message("test message str".getBytes());
set = new MessageSet(message);
producer.send("topic1", set);
若要订阅主题,消费者首先为该主题创建一个或多个消息流。发布到该主题的消息会被均匀分发到这些子流中,Kafka分发消息的具体细节将在3.2节介绍。每个消息流都提供了迭代器接口,用于遍历持续产生的消息流。消费者遍历流中的每条消息,处理消息的有效载荷。与传统迭代器不同,消息流迭代器永远不会终止。如果当前没有更多消息可供消费,迭代器会阻塞,直到有新的消息发布到该主题。我们既支持点对点交付模型(多个消费者共同消费主题中所有消息的单个副本),也支持发布-订阅模型(多个消费者各自获取主题的完整副本)。
消费者示例代码:
streams[] = Consumer.createMessageStreams("topic1", 1);
for (message : streams[0]) {
bytes = message.payload();
// do something with the bytes
}
Kafka的整体架构如图1所示。由于Kafka本身是分布式的,一个Kafka集群通常包含多个代理节点。为了实现负载均衡,一个主题会被划分为多个分区,每个代理节点存储其中一个或多个分区。多个生产者和消费者可以同时发布和拉取消息。在3.1节中,我们将介绍单个分区在代理节点上的存储布局,以及为提升分区访问效率所做的几项设计选择;3.2节介绍分布式场景下生产者、消费者与多个代理节点的交互方式;3.3节讨论Kafka的交付保证。

图1 Kafka架构
3.1 单分区效率
我们在Kafka中做了多项设计决策,以提升系统效率。
简洁的存储结构:Kafka采用非常简单的存储布局。主题的每个分区对应一个逻辑日志;物理上,一个日志由一组大小大致相同的段文件(例如1GB)实现。每当生产者向分区发布一条消息,代理节点只需将消息追加到最后一个段文件的末尾。为了提升性能,只有当发布的消息数量达到配置阈值,或者经过了指定的时间后,才会将段文件刷写到磁盘中。消息只有在被刷写之后,才对消费者可见。
与典型的消息系统不同,Kafka中存储的消息没有显式的消息ID。取而代之的是,每条消息通过其在日志中的逻辑偏移量来寻址。这种设计避免了维护辅助的、高随机访问开销的索引结构(用于映射消息ID到实际消息位置)的开销。需要注意的是,我们的消息ID是递增的,但不是连续的。要计算下一条消息的ID,需要将当前消息的长度加上其ID。下文我们将消息ID和偏移量作为等价概念使用。
消费者总是按顺序消费特定分区中的消息。如果消费者确认了某个消息偏移量,就意味着该消费者已经接收了分区中该偏移量之前的所有消息。在底层,消费者会向代理节点发送异步拉取请求,预先缓存一批数据,供上层应用消费。每个拉取请求都包含消费起始位置的消息偏移量,以及可接受的拉取字节数。每个代理节点都会在内存中维护一个排序的偏移量列表,包含每个段文件第一条消息的偏移量。代理节点通过检索偏移量列表,找到请求消息所在的段文件,再将数据返回给消费者。消费者收到消息后,计算出下一条待消费消息的偏移量,用于下一次拉取请求。Kafka日志的存储布局与内存索引结构如图2所示,每个方框标注了对应消息的偏移量。

图2 Kafka日志结构
高效的数据传输:我们对Kafka的数据传入与传出过程做了精细优化。前文提到,生产者可以在单次发送请求中提交一组消息。尽管上层消费者API是逐条遍历消息,但在底层,消费者的每次拉取请求也会获取多条消息,总大小通常可达数百KB。
我们做出的另一项非常规设计选择,是不在Kafka层面对消息进行显式的内存缓存,而是依赖底层文件系统的页缓存(page cache)。这种设计的核心优势是避免了双重缓冲——消息仅在页缓存中缓存一份。此外,即使代理节点进程重启,页缓存中的热数据依然有效。由于Kafka完全不在进程内缓存消息,其内存垃圾回收的开销极低,这使得使用基于虚拟机的语言实现高效系统成为可能。最后,由于生产者和消费者都是顺序访问段文件,且消费者通常仅落后生产者一小段距离,操作系统的常规缓存策略(特别是写通缓存与预读机制)会非常有效。我们发现,即使数据量达到数TB,生产与消费的性能依然与数据规模呈线性关系,表现稳定。
此外,我们还针对消费者的网络访问做了优化。Kafka是一个多订阅者系统,同一条消息可能会被不同的消费者应用多次消费。将本地文件中的数据发送到远程套接字的典型流程包含以下步骤:(1)从存储介质读取数据到操作系统的页缓存;(2)将页缓存中的数据复制到应用层缓冲区;(3)将应用层缓冲区复制到内核缓冲区;(4)将内核缓冲区的数据发送到套接字。这个过程包含4次数据拷贝和2次系统调用。在Linux及其他类Unix操作系统中,提供了sendfile API[5],可以直接将数据从文件通道传输到套接字通道,通常可以省去步骤(2)和(3)中的2次拷贝与1次系统调用。Kafka利用sendfile API,实现了将日志段文件中的数据从代理节点高效传输到消费者。
无状态的代理节点:与大多数其他消息系统不同,在Kafka中,每个消费者的消费进度信息不由代理节点维护,而是由消费者自身维护。这种设计大幅降低了代理节点的复杂度与开销。然而,这也使得消息删除变得棘手,因为代理节点无法知晓是否所有订阅者都已消费了某条消息。Kafka通过基于时间的简单服务等级协议(SLA)实现留存策略来解决这个问题:如果消息在代理节点中的留存时间超过指定时长(通常为7天),就会被自动删除。该方案在实践中效果良好,包括离线消费者在内的绝大多数消费者,都会按天、按小时或实时完成消费。而Kafka的性能不会随数据量增大而下降,也让这种长周期留存策略具备可行性。
这种设计还有一个重要的附加优势:消费者可以主动回溯到旧的偏移量,重新消费数据。这违反了队列的常规约定,但对于很多消费者而言却是一项核心特性。例如,当消费者的应用逻辑出现错误时,修复错误后可以重放特定的消息。这对于向数据仓库或Hadoop系统进行ETL数据加载的场景尤为重要。再比如,消费后的数据可能只会周期性地刷写到持久化存储中(例如全文索引器)。如果消费者崩溃,未刷写的数据就会丢失。这种情况下,消费者可以为未刷写消息的最小偏移量设置检查点,重启后从该偏移量重新开始消费。需要说明的是,相比于推送模型,拉取模型实现消费者回溯要容易得多。
3.2 分布式协调
接下来介绍分布式场景下生产者与消费者的工作机制。生产者发布消息时,可以选择随机分发到某个分区,也可以通过分区键与分区函数,按照业务语义分发到指定分区。本节将重点讨论消费者与代理节点的交互方式。
Kafka引入了**消费者组(consumer group)**的概念。每个消费者组由一个或多个消费者组成,共同消费一组订阅的主题——即每条消息只会交付给组内的一个消费者。不同的消费者组各自独立消费完整的订阅消息集合,组间不需要任何协调。同一消费者组内的消费者可以运行在不同进程中,也可以部署在不同机器上。我们的目标是,在不引入过多协调开销的前提下,将代理节点中存储的消息均匀分发给各个消费者。
我们的第一个设计决策是:将主题内的分区作为并行处理的最小单元。这意味着在任意时刻,每个分区的所有消息,只会被每个消费者组内的一个消费者消费。如果允许多个消费者同时消费同一个分区,它们就需要协调各自消费的消息范围,这会带来锁机制与状态维护的开销。相比之下,在我们的设计中,消费进程仅在消费者进行负载重平衡时才需要协调,而重平衡是低频事件。为了实现真正的负载均衡,我们要求主题的分区数量远多于每个消费者组内的消费者数量。通过对主题进行过度分区,我们可以轻松实现这一点。
我们的第二个设计决策是:不设置中心“主”节点,而是让消费者以去中心化的方式自行协调。引入主节点会增加系统复杂度,因为我们还需要处理主节点故障的问题。为了实现协调,我们采用了高可用的一致性服务ZooKeeper[10]。ZooKeeper提供了类似文件系统的简洁API,用户可以创建路径、设置路径的值、读取路径的值、删除路径,以及列出路径的子节点。它还具备几个更实用的特性:(a)用户可以在路径上注册监听器(watcher),当路径的子节点或路径的值发生变化时,监听器会收到通知;(b)路径可以被创建为临时节点(ephemeral,与持久节点相对),即如果创建该节点的客户端断开连接,ZooKeeper服务器会自动删除该路径;(c)ZooKeeper会将数据复制到多台服务器上,保证数据的高可靠性与高可用性。
Kafka将ZooKeeper用于以下任务:(1)检测代理节点与消费者的新增和移除;(2)当上述事件发生时,触发每个消费者的重平衡流程;(3)维护消费归属关系,并记录每个分区的消费偏移量。具体来说,每个代理节点或消费者启动时,都会将自身信息存储到ZooKeeper中的代理节点注册表或消费者注册表中。代理节点注册表包含代理节点的主机名、端口,以及其上存储的主题与分区集合;消费者注册表包含消费者所属的消费者组,以及其订阅的主题集合。每个消费者组在ZooKeeper中对应一个所有权注册表和一个偏移量注册表:所有权注册表为每个订阅分区维护一个路径,路径的值为当前消费该分区的消费者ID(我们称之为消费者“拥有”该分区);偏移量注册表存储每个订阅分区中最后一条已消费消息的偏移量。
ZooKeeper中创建的路径,代理节点注册表、消费者注册表和所有权注册表对应的是临时节点,偏移量注册表对应的是持久节点。如果某个代理节点故障,它上面的所有分区都会自动从代理节点注册表中移除;如果消费者故障,它在消费者注册表中的条目,以及它在所有权注册表中拥有的所有分区都会被清除。每个消费者都会在代理节点注册表和消费者注册表上注册ZooKeeper监听器,当代理节点集合或消费者组发生变化时,都会收到通知。
在消费者初始启动时,或者当消费者通过监听器感知到代理节点/消费者变更时,会发起重平衡流程,重新确定自己应当消费的分区子集。该流程如算法1所示。消费者首先从ZooKeeper读取代理节点注册表与消费者注册表,计算出每个订阅主题T的可用分区集合$P_T$,以及订阅该主题的消费者集合$C_T$。然后将$P_T$按范围划分为$|C_T|$块,并确定性地选择其中一块作为自己的分区。对于每个分配到的分区,消费者会在所有权注册表中将自己标记为该分区的新所有者。最后,消费者启动线程,从偏移量注册表中记录的偏移量开始,拉取每个所属分区的数据。随着分区消息的不断拉取,消费者会周期性地更新偏移量注册表中记录的最新消费偏移量。
算法1:消费者组G中消费者$C_i$的重平衡流程
For each topic T that Ci subscribes to {
remove partitions owned by Ci from the ownership registry
read the broker and the consumer registries from Zookeeper
compute PT =partitions available in all brokers under topic T
compute CT =all consumers in G that subscribe to topic T
sort PT and CT
let j be the index position of Ci in CT and let N =|PT|/|CT|
assign partitions from j*N to (j+1)*N -1 in PT to consumer Ci
for each assigned partition p {
set the owner of p to Ci in the ownership registry
let Op =the offset of partition p stored in the offset registry
invoke a thread to pull data in partition p from offset Op
}
}
当组内存在多个消费者时,所有消费者都会收到代理节点或消费者变更的通知,但通知到达各个消费者的时间可能略有差异。因此,可能出现一个消费者尝试获取另一个消费者仍拥有的分区所有权的情况。发生这种情况时,前者会释放自己当前拥有的所有分区,等待一段时间后重试重平衡流程。在实践中,重平衡流程通常仅需重试几次就会稳定。
当创建一个新的消费者组时,偏移量注册表中没有对应的偏移量记录。此时,消费者会通过代理节点提供的API,从每个订阅分区的最小偏移量或最大偏移量开始消费(具体取决于配置)。
3.3 交付保证
总体而言,Kafka仅保证至少一次交付。精确一次交付通常需要两阶段提交,对于我们的应用场景而言并非必要。在大多数情况下,消息会精确一次地交付给每个消费者组。然而,如果消费者进程在未正常关闭的情况下崩溃,接管故障消费者所属分区的新消费者进程,可能会收到部分重复的消息——这些消息的偏移量位于最后一次成功提交到ZooKeeper的偏移量之后。如果应用不能容忍重复消息,必须自行实现去重逻辑,既可以利用返回给消费者的偏移量,也可以利用消息内的唯一键。相比于使用两阶段提交,这种方案通常具备更高的成本效益。
Kafka保证,单个分区内的消息会按顺序交付给消费者。但是,对于来自不同分区的消息,Kafka不保证其整体顺序。
为了避免日志损坏,Kafka在日志中为每条消息都存储了CRC校验值。如果代理节点发生I/O错误,Kafka会运行恢复流程,移除CRC校验不一致的消息。消息级别的CRC校验,也让我们可以在消息生产或消费后检查网络传输错误。
如果某个代理节点宕机,存储在该节点上尚未消费的消息将无法访问。如果代理节点的存储系统发生永久性损坏,所有未消费的消息都会永久丢失。未来,我们计划在Kafka中内置复制机制,将每条消息冗余存储到多个代理节点上。
4. 领英的Kafka应用实践
本节介绍领英对Kafka的使用方式。图3展示了我们部署架构的简化版本。我们在每个部署用户-facing服务的数据中心,都配套部署了一个Kafka集群。前端服务生成各类日志数据,并批量发布到本地的Kafka代理节点。我们通过硬件负载均衡器,将发布请求均匀分发到Kafka代理节点集群。Kafka的在线消费者运行在同一数据中心的服务中。

图3 Kafka部署架构
我们还在一个独立的数据中心部署了Kafka集群,用于离线分析,该数据中心在地理位置上靠近我们的Hadoop集群与其他数据仓库基础设施。这个Kafka实例运行了一组嵌入式消费者,从各个线上数据中心的Kafka实例中拉取数据。之后,我们运行数据加载作业,将这个Kafka副本集群中的数据拉取到Hadoop和数据仓库中,在其上运行各类报表任务与分析流程。我们也利用这个Kafka集群进行原型验证,还可以针对原始事件流运行简单脚本,实现即席查询。在未做大量调优的情况下,整个端到端管道的平均延迟约为10秒,完全满足我们的需求。
目前,Kafka每天累计处理数百GB数据、近十亿条消息;随着我们逐步将遗留系统迁移到Kafka,预计数据量还会显著增长,未来还会接入更多类型的消息。当运维人员因软件或硬件维护而启动或停止代理节点时,重平衡流程能够自动重新分配消费任务。
我们的监控体系还包含一套审计系统,用于验证整个数据管道中不存在数据丢失。为此,每条消息都携带了生成时的时间戳与服务器名称。我们为每个生产者都做了埋点,使其周期性地生成监控事件,记录固定时间窗口内该生产者针对每个主题发布的消息数量。生产者将监控事件发布到Kafka的一个独立主题中。消费者可以统计自己从指定主题接收的消息数量,并与监控事件中的计数进行比对,从而验证数据的正确性。
向Hadoop集群加载数据的功能,是通过实现一个专用的Kafka输入格式完成的,它允许MapReduce作业直接从Kafka读取数据。MapReduce作业加载原始数据后,会对数据进行分组与压缩,以便后续高效处理。无状态代理节点与客户端侧维护消息偏移量的设计,在此处再次发挥了作用——MapReduce的任务管理机制(支持任务失败与重启)可以以天然的方式处理数据加载,在任务重启时不会出现消息重复或丢失。只有作业成功完成时,数据和偏移量才会被存储到HDFS中。
我们选择Avro[2]作为序列化协议,因为它效率高且支持模式演进。对于每条消息,我们都会在有效载荷中存储其Avro模式ID与序列化后的字节数据。该模式让我们可以强制约束契约,确保数据生产者与消费者之间的兼容性。我们使用一个轻量级的模式注册服务,将模式ID映射到对应的实际模式。消费者收到消息后,会在模式注册表中查找对应的模式,将字节数据解码为对象(由于模式是不可变的,每个模式只需查找一次)。
5. 实验结果
我们开展了一项实验研究,将Kafka与两款主流消息系统进行性能对比:Apache ActiveMQ v5.4[1]——一款流行的JMS开源实现,以及RabbitMQ v2.4[16]——一款以性能著称的消息系统。我们使用ActiveMQ默认的持久化消息存储KahaDB。我们还测试了ActiveMQ的另一种消息存储,其性能与KahaDB非常接近,此处不再展示。在所有系统中,我们都尽可能使用可对比的配置参数。
实验在2台Linux机器上运行,每台机器配置8核2GHz处理器、16GB内存、6块磁盘组成RAID 10阵列,两台机器通过1Gb网络连接。其中一台机器作为代理节点,另一台作为生产者或消费者。
生产者测试:我们将所有系统的代理节点配置为异步将消息刷写到持久化存储。针对每个系统,我们运行单个生产者,总共发布1000万条消息,每条消息大小为200字节。我们将Kafka生产者的批量大小分别配置为1和50。ActiveMQ和RabbitMQ没有便捷的消息批量发送方式,我们默认其批量大小为1。测试结果如图4所示,横轴为随时间推移发送到代理节点的数据量(单位:MB),纵轴对应生产者吞吐量(单位:消息/秒)。平均来看,当批量大小为1和50时,Kafka的消息发布速率分别为5000条/秒和40000条/秒。这一数值比ActiveMQ高出数个数量级,至少是RabbitMQ的2倍。

图4 生产者性能
Kafka性能更优的原因主要有几点。首先,Kafka生产者当前不需要等待代理节点的确认,而是以代理节点能处理的最快速度发送消息,这显著提升了生产者的吞吐量。当批量大小为50时,单个Kafka生产者几乎可以打满生产者与代理节点之间的1Gb链路。对于日志聚合场景而言,这是一种合理的优化,因为数据需要异步发送,避免给线上业务流量引入任何延迟。需要注意的是,生产者不等待确认的机制,无法保证每条发布的消息都能被代理节点实际接收。对于很多类型的日志数据,只要丢失的消息数量相对较少,用持久性换取吞吐量是可接受的。不过,我们计划在未来针对更关键的数据场景,解决持久性保障的问题。
其次,Kafka的存储格式更高效。平均而言,Kafka中每条消息的额外开销仅为9字节,而ActiveMQ为144字节。这意味着存储相同的1000万条消息,ActiveMQ比Kafka多占用70%的空间。ActiveMQ的开销一方面来自JMS规范要求的厚重消息头,另一方面来自维护各类索引结构的成本。我们观察到,ActiveMQ中最繁忙的线程之一,大部分时间都在访问B树以维护消息元数据与状态。最后,批量发送通过分摊RPC开销,极大地提升了吞吐量。在Kafka中,50条消息的批量大小,让吞吐量提升了近一个数量级。
消费者测试:第二个实验测试了消费者的性能。同样,针对所有系统,我们使用单个消费者总共消费1000万条消息。我们将所有系统的预取大小配置为相近水平——每次拉取请求最多获取1000条消息,约合200KB。对于ActiveMQ和RabbitMQ,我们将消费者确认模式设置为自动确认。由于所有消息都可以存放在内存中,所有系统都从底层文件系统的页缓存或内存缓冲区中提供数据。测试结果如图5所示。
平均来看,Kafka的消费速度为22000条消息/秒,是ActiveMQ与RabbitMQ的4倍以上。我们认为有几方面原因:首先,Kafka的存储格式更高效,因此从代理节点传输到消费者的数据量更少;其次,ActiveMQ与RabbitMQ的代理节点都需要维护每条消息的交付状态。我们观察到,测试过程中ActiveMQ的一个线程一直在忙于将KahaDB页面写入磁盘。相比之下,Kafka代理节点没有任何磁盘写入活动。最后,通过使用sendfile API,Kafka降低了传输开销。

图5 消费者性能
在本节最后需要说明:本实验的目的并非证明其他消息系统不如Kafka。毕竟,ActiveMQ和RabbitMQ都具备比Kafka更丰富的功能。本实验的核心意义在于,说明专用系统可以实现的性能提升潜力。
6. 结论与未来工作
本文提出了一款名为Kafka的新型系统,用于处理海量日志数据流。与消息系统类似,Kafka采用基于拉取的消费模型,允许应用按照自身速率消费数据,并可在需要时回溯消费。通过聚焦日志处理场景,Kafka实现了比传统消息系统高得多的吞吐量,同时提供了完善的分布式支持,具备水平扩展能力。我们已在领英成功将Kafka应用于离线与在线两类场景。
未来,我们有多个方向的工作规划。首先,我们计划在Kafka中内置跨代理节点的消息复制机制,即使发生不可恢复的机器故障,也能保证数据持久性与可用性。我们将同时支持异步与同步复制模型,允许用户在生产者延迟与保障强度之间做权衡。应用可以根据自身对持久性、可用性与吞吐量的要求,选择合适的冗余级别。其次,我们希望在Kafka中增加流处理能力。实时应用从Kafka获取消息后,通常会执行类似的操作,例如基于窗口的计数、将每条消息与二级存储中的记录或其他流中的消息进行关联。在最底层,我们可以通过在发布时基于关联键对消息进行语义分区来支持这一点——所有携带特定键的消息都会进入同一个分区,从而被同一个消费者进程处理。这为跨消费者集群处理分布式流提供了基础。在此之上,我们认为提供一套实用的流处理工具库(例如各类窗口函数、关联技术),将对这类应用非常有价值。
7. 参考文献
[1] http://activemq.apache.org/
[2] http://avro.apache.org/
[3] Cloudera’s Flume, https://github.com/cloudera/flume
[4] http://developer.yahoo.com/blogs/hadoop/posts/2010/06/enabling_hadoop_batch_processi_1/
[5] Efficient data transfer through zero copy: https://www.ibm.com/developerworks/linux/library/j-zerocopy/
[6] Facebook’s Scribe,
http://www.facebook.com/note.php?note_id=32008268919
[7] IBM Websphere MQ: http://www01.ibm.com/software/integration/wmq/
[8] http://hadoop.apache.org/
[9] http://hadoop.apache.org/hdfs/
[10] http://hadoop.apache.org/zookeeper/
[11] http://www.slideshare.net/cloudera/hw09-hadoop-based-data-mining-platform-for-the-telecom-industry
[12] http://www.slideshare.net/prasadc/hive-percona-2009
[13] https://issues.apache.org/jira/browse/ZOOKEEPER-775
[14] JAVA Message Service: http://download.oracle.com/javaee/1.3/jms/tutorial/1_3_1-fcs/doc/jms_tutorialTOC.html.
[15] Oracle Enterprise Messaging Service: http://www.oracle.com/technetwork/middleware/ias/index-093455.html
[16] http://www.rabbitmq.com/
[17] TIBCO Enterprise Message Service: http://www.tibco.com/products/soa/messaging/
[18] Kafka, http://sna-projects.com/kafka/