上次更新时间:06/12/2026
Apache Kafka 是一个热门的开源事件流处理平台,用于收集、处理和存储连续的事件流。Kafka 通常用作消息传递中间件,但它具有可伸缩性和冗余性,使分布式应用能够处理从每天一个事件到每秒数十亿个事件的各种情况。与传统的消息传递系统不同,Kafka 也是一个持久存储系统,它将记录存储在有序日志中,可以重复读取。这使得 Kafka 成为向事务数据库分发更改的常用系统,可用于在分析和其他系统中重建或实现数据。这种模式有时称为事件溯源。
因此,Kafka 对于传统事件总线模式(其中事件驱动型应用通过消息传递中间件集成)和数据联合(或“异构具体化”)架构都很重要。这种方法延迟低且经济高效。
探索 Google Cloud Managed Service for Apache Kafka,自动执行流式处理基础设施,并加快“数据到 AI”工作流的速度。
Kafka 采用流式数据,能够准确记录何时发生了什么。此记录称为仅允许附加操作的日志。它不可改变,因为它可以被附加,但不能改变。然后,应用可以订阅日志以访问数据,也可以向日志发布数据以实时添加更多数据。
“Kafka”的核心通常是指低延迟存储系统,但流式处理平台还包括其他重要组件。首先是 Kafka Connect,这是一个集成系统,允许将可横向扩缩的连接器连接到许多重要系统。这包括变更数据捕获 (CDC) 连接器、集群到集群复制 (MirrorMaker),以及将数据写入下游系统(如湖仓一体 [Apache Iceberg]、数据湖 [对象存储上的 Avro 或 Parquet 文件])以及 BigQuery 等数据库的功能。其次,Kafka 项目附带了一组强大的客户端,包括用于操作集群和主题的管理命令行客户端,以及用于读写数据的高性能客户端库。
过去,数据处理是通过周期性批量作业进行的,即先存储原始数据,然后以任意时间间隔进行处理。例如,零售公司可能会等到一天结束时才分析销售数据。批处理的局限性之一在于它不是实时的。越来越多的组织和数据科学家希望在数据生成时实时分析数据,以便及时做出业务决策并为实时 AI 模型提供支持。这就是事件流的用武之地。事件流是连续不断地处理无限数据流(自其创建之时起)的过程。这可以捕获数据的时间价值,并支持基于推送的应用,以便在重要事情发生时及时采取行动。对于数据科学家来说,这意味着能够执行实时特征工程并提供低延迟预测。
虽然许多组织都专注于数据科学家生成的下游分析洞见,但 Apache Kafka 的主要从业者是数据工程师。这些专业人员负责构建连接公司应用和数据库的关键“数据管道”和集成。
数据工程师使用 Kafka 在整个技术栈中创建可靠的连接。这些集成可以采用多种形式:
在典型的企业中,数据工程师与应用团队密切合作,确保将用户事件、业务交易和数据库更新导出到 Kafka。这一过程称为数据联合,可让组织中的多个用户和系统同时使用这些事件。
在数据科学和 AI 领域,Kafka 的价值主要体现在数据访问方面。它可作为应用日志和数据库变更的全面来源。Kafka 以速度快而闻名,但对于大多数数据科学工作流而言,数据源的广度和可靠性远比低延迟交付更重要。
AI 系统在推理期间依赖高质量的训练数据和上下文。Kafka 通常对于训练至关重要,它可以从各种源系统收集训练数据,包括互动日志和数据库更改。在许多组织中,它被用作事件总线,用于聚合来自许多服务的事件,或者只是用作应用日志的暂存位置。这使其成为生成训练数据集的天然单一数据源。由于 Kafka 以有序序列存储记录,因此它也特别适合处理序列的 LLM。
Kafka 对于许多在线推理任务至关重要。应用或智能体能否提供相关的产品推荐、搜索回答或提示,取决于是否掌握了用户的最新情境。由于 Kafka 支持低延迟、可扩缩的通信,因此推理系统可以在几十毫秒内使用最新事件更新用户上下文。例如,如果用户拒绝了音乐应用中最新推荐的歌曲,或者金融应用中的股票价格发生了变化,推荐服务可以立即生成更好的建议,并将这些输入考虑在内。
开源生态系统
Kafka 的源代码是免费提供的,得益于全球社区,它拥有各种连接器、监控工具和插件。
规模和速度
Kafka 是一个分布式平台,这意味着处理过程被分配给多台机器。这使其能够扩缩以处理海量数据,同时保持亚毫秒级延迟。
高可用性
由于 Kafka 是分布式系统,即使单个机器发生故障,它也能保持可靠性,因此适合用于关键任务应用
设置本地 Kafka 集群非常困难,需要团队预配机器、管理安全并处理日常修补。借助托管式服务,提供商会处理底层基础设施,让您可以专注于构建应用。这对于想要专注于模型开发和分析洞见而非基础设施管理的数据科学团队来说尤其有益。
Kafka 通过以下四个核心功能实现流式事件处理: