WRITING / 2026.10.02

大数据专题(二):Kafka 数据接入,分区、消费进度与可靠重放

从视频播放事件的接入与积压出发,理解分区键、消费位置、提交顺序和幂等边界,用真实 Kafka 实验观察恢复消费与主动回放。

一次重启为什么可能让播放量上涨

假设消费程序读到一个播放事件,把它追加到明细文件,然后提交 Kafka 消费位置。正常运行时,这两个动作紧挨着发生,很容易被当成一个整体。可是进程可能恰好在文件写完、位置提交之前退出。重启后,消费者从上一次已提交的位置继续读取,刚才那条事件又出现一次。

如果下游直接对每一行执行加一,播放量就多算了。事件没有被 Kafka 凭空复制,问题出在两个系统之间的完成边界没有对齐。反过来,先提交位置再写文件,失败时又可能跳过尚未写出的事件。这是数据接入中最值得先弄清楚的故障窗口。

本篇不把“用了消息队列”当作可靠性的答案,而是跟着一条合成播放事件观察:它被写到哪里、何时算接受成功、消费者如何记住进度,以及为什么允许重放之后还需要业务去重。

Kafka 分区日志、消费进度与重放的关系

日志保留与消费进度是两个问题

Kafka topic 可以划分为多个分区,每个分区保存有序的记录序列。记录在分区中的位置由 offset 标识。对于读者,最有用的理解方式是把“日志里存在什么”和“某个消费者读到哪里”分开:同一份事件可以被日报任务、实时榜单和审计程序分别读取,各自推进进度。

一组消费者读完之后,并不会因此立即删除记录。记录是否继续保留,由 topic 的保留或压缩策略等配置决定。也正因为保留与读取解耦,消费者才有机会重新读取历史范围。具体日志、消费者位置和复制语义可参考 Kafka 3.9 设计文档。

这种结构让消费程序能够停下来追赶,也引入一个现实约束:它必须在需要的数据被清理之前恢复。把保留时间设置为一天,而允许日报系统连续停止三天,是互相矛盾的恢复目标。保留周期应结合最大可接受停机时间、补数范围和磁盘预算确定。

对于已经超出保留范围的历史,不能靠重置 offset 把数据找回来。此时必须依赖另一个保留原始数据的存储,或重新从权威源导入。Kafka 可以承担事件缓冲和有限历史重放,但是否足以作为唯一历史来源,需要按业务恢复要求判断。

事件身份与 offset 为什么不能互换

event_id 标识业务事件,topic + partition + offset 标识该日志中的一次记录位置。一条业务事件可以因为生产端重试之外的应用行为、手动回放或上游重复上报,出现在多个不同 offset 上。仅按 offset 去重,可以避免同一条日志记录被重复应用,却无法识别不同记录实际代表同一事件。

反过来,两个真正不同的播放事件即使字段完全相同,也不能因为用户、视频和观看时长一样就合并。重复播放是正常行为。事件 ID 应在业务事件形成时生成,重传时保持稳定,而不是每次向 Kafka 发送前临时生成一个新 ID。

本实验保留这两个身份:Kafka 管理读取位置,计算器使用事件 ID 识别业务重复。确定性样例会将一百条事件原样再发一次,产生新的日志记录但沿用原 ID。最终合法唯一事件仍应是一万条。这个测试比只制造网络重试更直接地检验业务去重。

同一 ID 对应不同业务内容则是另外一种错误。生产系统需要定义拒绝、隔离或按版本更新的政策,不能默默选择一个并宣称数据已去重。本系列基础样例仅包含同 ID 同内容的重复,独立参考计算器遇到冲突会报错;这条边界需要在扩展真实业务前重新设计。

分区键决定哪些记录能够保持相对顺序

本实验生产者按视频 ID 选择分区,使同一视频的记录进入同一分区。在分区集合固定、分区器行为一致的条件下,这便于观察某个视频的事件顺序与消费位置。它并不提供整个 topic 的全局顺序:不同分区可以由不同任务并行读取,到达下游的先后关系也可能变化。

选择键时必须先确认真正需要顺序的对象。如果处理的是订单状态转换,订单 ID 可能更合适;如果统计视频热度,按视频分区又可能让热门视频成为热点。一种键既满足顺序又满足均匀分布,往往需要业务条件配合,不能只依靠一个通用哈希函数。

增加分区可以提高可用并行度,但也可能改变键到分区的映射。于是,原本位于旧分区的某个键,其后续记录可能进入新分区。对于依赖同键连续处理的系统,扩容必须同时考虑顺序与状态迁移,不能把“分区数增大”理解为纯粹的容量按钮。

本系列在接入阶段使用三个固定分区,后续 Flink 再按事件 ID 去重、按视频 ID 聚合。输入的物理分区和计算算子的业务分组是两层不同组织方式。理解这一点,有助于解释为什么一个看似已经按视频分过区的输入,计算过程中仍可能发生数据交换。

消费组、当前位置与已提交位置

使用消费组订阅时,同一组中的成员协作分配分区,不同消费组则保留独立进度。成员加入、退出或订阅变化会触发分配调整。组内并行消费的上限和分区数量有关,但单个实例的吞吐也受处理时间、批量大小与下游写入影响。

消费者当前读取位置表示接下来准备读取哪里;已提交位置表示应用希望在恢复时从哪里继续。这两个位置可能不同。通常提交的是下一条待处理记录的位置,而不是最后一条已完成记录的位置。自己保存 offset 时应明确这项约定,否则容易产生一条数据的偏移错误。

为了将实验聚焦到进度语义,配套 KafkaTool 使用显式分区分配,捕获每个分区的结束位置,并执行有界导出。它没有演示动态组成员再均衡,因此不能用这个小程序的结果评价再均衡稳定性。生产系统采用订阅模式时,还需要处理分区撤销期间的状态和未完成写入。

有界导出也不能靠“连续几秒没有消息”判断完成。网络抖动、事务未提交或分区短暂无数据都可能造成空轮询。本实验用开始时记录的结束 offset 作为边界,直到各分区读取位置达到它;超时明确失败,不将空文件伪装为成功。

提交顺序背后的交付语义

先提交消费位置,再执行外部写入,意味着恢复后可能不再读取已经确认但尚未落地的记录。先完成外部写入,再提交位置,意味着故障后可能重复应用已经落地的记录。前者偏向避免重复但可能遗漏,后者偏向避免遗漏但必须能处理重复。

这里的“完成写入”还要继续追问。写进语言运行时缓冲区、写进操作系统缓存、完成文件关闭、远端返回成功,不一定代表相同的持久性。示例导出工具用于构造实验输入,并不是生产级文件归档器;它没有实现文件清单和消费位置的跨系统原子提交。

工程上可以选择让目标写入幂等:使用稳定业务键覆盖最终值,或在同一数据库事务中记录已处理 ID 与业务结果。也可以让输出和进度进入一个共同事务边界。哪一种方式合适,取决于目标系统支持什么,而不是仅看消费者配置项。

Kafka 的幂等生产者处理的是特定生产协议范围内的重复发送问题;应用重新构造一条同业务 ID 的消息再次发送,仍可能形成新的记录。事务生产者进一步提供事务可见性,但不会自动使任意 HTTP 调用、文件追加或邮件发送具有同样的事务语义。生产者配置应与目标系统约束一起阅读。

Exactly-once 要写清楚起点与终点

本系列第五篇采用 Kafka 输入、Flink 状态和 Kafka 事务输出,消费者以 read_committed 读取已提交结果。这是一条能够明确描述输入、状态和输出边界的实验路径。它仍需要稳定的算子身份、正确的 Checkpoint、合适的事务超时以及唯一的事务 ID 前缀。

即便这条路径在故障实验中验证通过,也不表示源头事件一定完整。客户端没有发送的事件、被错误过滤的事件、两个不同 ID 表示同一业务行为的情况,都不在传输和状态一致性的自动保证范围内。交付语义负责某个已定义输入的处理效果,业务语义仍由应用负责。

观察结果时还要区分重复交付和合法更新。窗口第一次输出播放量二,迟到事件被接纳后输出播放量三,这是对同一窗口结果的修订,不是重复记录。下游如果把二和三再次相加,会制造五这个错误答案。数据接口必须说明值代表增量还是最新绝对值。

最小实验:积压、恢复与主动回放

下载源码包,按 README准备固定版本。Kafka 使用独立数据目录和回环端口,单节点同时承担 broker 与 controller,副本数为一。这个配置降低了学习成本,也意味着本实验没有验证副本故障容忍。

python3 run.py setup --work /absolute/path/to/lab-work
python3 run.py kafka --work /absolute/path/to/lab-work

接入步骤先在归档消费者未运行时写入全部样例,再启动消费者追赶。首次有界导出得到 10,101 条记录;使用相同消费组从已提交位置恢复,新增导出为零条;主动从最早位置回放到另一份文件,再次得到 10,101 条。另一个独立主题专门调用 pause():暂停后生产 10,101 条记录,poll() 返回零条且位置保持不动,记录积压 10,101;调用 resume() 后完整导出 10,101 条,去重参考一致。这验证显式暂停期间积压得以保留,但没有模拟消费组再均衡。把导出结果交给独立参考计算器后,得到 10,000 条唯一合法事件,与原始答案一致。

操作 本次实际导出行数 解释
首次读取积压 10,101 包含 100 条业务重复和 1 条非法样例
从已提交位置恢复 0 本次没有追加新的播放事件
主动从头回放 10,101 保留中的记录仍可重新读取

这三次读取没有使用“消费后删除消息”的模型。第二次读不到新增记录,并不说明日志已经为空;第三次能够回放,正好说明日志保留与组进度相互独立。实验日志记录各分区的结束位置,便于核对导出范围。

本次也没有声称导出器能在任意崩溃点实现精确一次归档。读者可以进一步在文件写入后、提交 offset 前注入退出,观察同一批数据再次出现;要将它变成可靠归档服务,应补充批次清单、目标提交与进度协调,而不是只给导出工具加一个无限重试循环。

积压为什么不能只看一个总数

消费落后可能是上游突发增长,也可能是某个分区倾斜、目标数据库变慢、反复重试或组成员变化。总积压相同的两次事故,恢复方式可能不同。排查时应同时观察每分区积压、最老未处理事件的年龄、消费速率、输入速率以及下游写入延迟。

offset 差值也不是精确的业务事件计数。事务控制记录、压缩后的空洞和过滤逻辑都可能让日志位置差与有效业务条数不同。它适合表达读取进度,不宜直接作为财务级或业务级完整性指标。业务对账仍应使用事件身份和明确的数据范围。

恢复速度需要大于持续输入速度,积压才会下降。增加消费者之前,应判断瓶颈是否可以并行化。如果所有记录都落到一个热点键,或所有任务最终竞争同一把数据库锁,增加实例可能只增加连接数。接入层的容量设计必须连同后续计算和存储一起评估。

埋点与 CDC 解决不同来源问题

播放事件来自客户端或服务端埋点,记录的是一次业务观察;CDC 捕获数据库变化,记录的是持久化状态的变更。前者适合表达“发生了什么行为”,后者适合将订单、用户、视频元数据等表的变化同步到分析系统。二者可能关联,却不能随意互相代替。

例如,视频标题修改可以通过 CDC 更新维度,而播放时长通常不能从视频元数据表的变化中还原。CDC 还要处理初始快照与增量衔接、更新与删除、源端日志保留、表结构变化和消费位点。完整接入应明确这些阶段的边界;本篇只比较职责,没有部署 CDC 连接器或宣称验证了源库快照一致性。

如果事务中写业务表,事务后单独发送消息,进程可能在两个动作之间失败。事务 Outbox 等模式尝试把待发送事实先与业务变化一起保存,再独立传递;下游仍需要识别重复。这类设计值得在真实业务接入时单独评审,不能把示例中的随机数据生产者直接替换成线上数据库写入代码。

重放要有独立的操作记录

主动回放最好记录发起原因、输入起止位置、目标位置和口径版本。教学实验将结果写到另一份文件,就是为了保留原结果作为比较依据。真实链路直接重置线上消费组位置,可能影响同组所有成员,还可能让正在展示的结果突然改变;应先在隔离目标中验证,再按明确规则替换或合并。

重放结束后还要检查两件事:这次是否覆盖了约定范围,目标是否仍满足幂等规则。如果只看到消费者重新运行而没有结果对账,就不能说明一次补数已经完成。控制面操作成功与业务修复成功,需要各自的证据。

从可靠接入走向可查询数据

Kafka 能让事件被保存、分区和重放,但不会把它们自动组织成一个便于分析、能够演进的表。下一步需要决定原始记录如何保留,哪些错误进入隔离区,如何生成去重明细,以及查询端怎样找到本次有效的数据文件。

把这些职责写清楚之后,接入层才有一个可验证的完成标准:能描述输入范围,能定位消费进度,能解释重复和缺失,能把合法数据交给下一层。第三篇将从导出的同一批事件出发,观察文件格式与分析表分别解决什么问题。

大数据工程:从业务事件到可信指标

  1. 大数据专题(一):从业务数据库到数据平台,一条分析链路如何形成
  2. 大数据专题(二):Kafka 数据接入,分区、消费进度与可靠重放 · 当前文章
  3. 大数据专题(三):从文件到分析表,Parquet、数仓分层与 Iceberg
  4. 大数据专题(四):Spark 批处理,一条 SQL 如何变成分布式计算
  5. 大数据专题(五):Flink 实时计算,乱序事件如何得到正确结果
  6. 大数据专题(六):让指标可信,批流对账、历史补数与故障恢复

按时间浏览