开发者生态
morning
超越偏移量延迟:PB 级规模下 Hudi 数据湖流水线的队列等待时间计算
2026-09-03
1 阅读
约10分钟阅读
作者:Srikanth Mamidala
字号:
问题背景 在 Twilio ",数据湖是公司各产品线(消息、邮件、语音等)进行数据分析、报表生成和机器学习的基础。内部团队通过这些数据来了解产品使用情况、驱动业务决策,并为异常检测或欺诈检测等机器学习模型提供支持。数据湖的流入管道使用 Apache Hudi Delta Streamer " 从 Kafka 摄取数据,截至 2025 年第四季度,每月处理超过 五万亿条记录 ",横跨自托管的 Kafka 集群,在 2025 年网络星期一 "这一天峰值达到每秒 1290 万条消息。在如此大的 规模 " 下,我们意识到必须给管道所有者提供一种精确的方式来定义和执行自定义的数据新鲜度 SLA(服务等级协议)。我们需要一种可操作的信号,同时不能给正在运行的管道增加任何额外开销。 传统的消费者延迟指标,如消费者偏移量延迟(records-lag-max)以及 Hudi 的 kafkaDelayCount,看起来都正常。消费者似乎跟得上 Kafka 的消息生产速度,但下游分析团队不断报告数据过期,有时甚至滞后数小时。问题不在于 Kafka 的吞吐量,而在于一个可见性盲区。Hudi Delta Streamer 管理自己的检查点,这些检查点与表数据一起存储在 S3 中,与 Kafka 的消费者组偏移量追踪是分离的。 Burrow " 这类标准延迟监控工具追踪的是消费者组提交的偏移量,而 Hudi 默认不会填充这些偏移量,因此它们无法感知 Hudi 是否已将数据实际提交到数据湖。 我们真正需要回答的问题是:最新的 Hudi 提交与当前 Kafka 主题中存在的消息之间到底落后了多少? 重新以时间维度看待延迟 Hudi 作业相对于 Kafka 主题中的消息到底落后了多久?换句话说,我们要计算自成功完成一次 Hudi 提交后,第一个未被消费的消息到达 Kafka 主题以来,已经过去了多长时间。 Hudi 中的偏移量追踪机制 HoodieStreamer "(前身为 HoodieDeltaStreamer)利用检查点机制来精确追踪哪些数据已被摄取,并防止重复处理相同的数据。对于 Kafka 数据源,检查点记录的是各个分区下已经成功处理并提交到存储的 Kafka 主题偏移量,或是时间戳。 检查点存储:检查点直接嵌入在 .hoodie 提交文件中,键名为 deltastreamer.checkpoint.key。 容错机制:在发生故障或重启时,HoodieStreamer 会读取最新的提交文件,获取该键值,并从上次中断的偏移量位置恢复读取。 Kafka 偏移量存储格式:以字符串形式存储,格式为 topicName,0:offset0,1:offset1。 根据这一事件时间线,我们可以使用 Apache Hudi SDK 找到最新的提交记录,并提取最后成功写入数据湖的各分区偏移量。 这种方法的实用之处在于,它不需要管道本身做任何新的事情。Hudi 已经提交到 S3 的偏移量和 Kafka 消息上已有的时间戳就足以计算出真实的数据新鲜度。我们通过一个叫作指标报告器的外部作业来计算新鲜度。指标报告器纯粹是一个外部观察者,读取系统已经产出的元数据,既不需要增加新的埋点,也不需要修改生产者端。 图 1. 采用指标报告器的高层数据管道示例。 数据摄入路径如下图 2 所示。Hudi Delta Streamer 将数据以及对应的偏移量检查点提交至 S3。指标上报器读取 Hudi 时间线,获取最新已提交的偏移量。 图 2. 指标报告器从 Hudi 时间线读取偏移量。 算法工作原理 在生产环境中,指标上报器每 15 分钟运行一次。针对每一条数据管道,该上报器会执行以下任务: 从 S3 上的活跃时间线中获取最新的 Hudi 提交记录。按照逆时间顺序遍历提交记录(从最新的开始),找到包含 deltastreamer.checkpoint.key 的最近一次提交。这样就可以拿到上一批成功提交数据对应的各分区偏移量。这就是 Hudi 已经消费完成并写入数据湖的偏移量位置。 在每个 Kafka 分区中定位到检查点对应的偏移量。这样可以将消费者定位到Hudi 尚未提交写入数据湖的第一条消息。这条消息仍然保留在 Kafka 中,是在上一次 Hudi 成功写入之后才到达的消息。 读取该消息,获取其时间戳 X。该时间戳代表数据进入 Kafka 的时间;这条数据一直驻留在 Kafka 中,等待下一轮 Hudi 任务处理。 计算延迟:当前时间戳 − X = 该数据已经等待的时长。 如果延迟超过阈值,则将上限设置为 7 天。该上限用于避免管道闲置或停止运行时出现无限大的指标数值。如果在检索范围内没有找到有效的检查点,则直接不输出该指标,避免上报具有误导性的错误数值。 快速示例 为便于理解这种方法,下面结合真实数值做完整演示。假设有一个叫作 orders‑events 的主题,包含 3 个分区。最新的 Hudi 提交包含以下检查点: orders-events,0:1200,1:980,2:1450 这些是下一次要读取的偏移量。Hudi 已提交分区 0 上偏移量 1199 及之前的所有消息,分区 1 上 979 及之前的所有消息,分区 2 上 1449 及之前的所有消息。报告器将每个分区定位到其检查点偏移量并拉取下一条记录。它从每个分区获取一条候选消息: 所有分区中最早的时间戳是四十五分钟前的。这就是当前等待被提交到数据湖的最旧消息。延迟为四十五分钟。如果该管道的 slaInMinutes 为三十分钟,则 SLA 比率为 45/30 = 1.5,上限为 1.0,表示完全违反 SLA。比率始终上限为 1.0,因为管道不可能超出“完全违反”。使用分级比率是一个设计选择,具体原因见“真实的结果”章节。 下面的代码将分三部分讲解实现逻辑:完整解决方案对应的核心算法;通过解析 Hudi 时间线(遍历至配置的最大检索深度)从 S3 获取 Hudi 检查点;以及定位到 Kafka 偏移量,读取尚未被消费的第一条消息。 核心算法(Java) 在下面的代码中,HoodieResult 是一个小的包装类,用于保存提交元数据和解析后的分区偏移量;kafkaClient 封装了一个标准的 Kafka Consumer;nextRecord 执行后文展示的查找和拉取操作。 // 步骤1:从S3读取Hudi时间线,提取检查点偏移量 HoodieResult hudiResult = findLatestCommitWithCheckpoint(tableName, tableBasePath, maxCommitDepth); final Map checkpointOffsets = hudiResult.getPartitionToCheckpoint() .entrySet().stream() .collect(Collectors.toMap( e -> new TopicPartition(topic, e.getKey()), Map.Entry::getValue )); // 步骤2、3:定位到检查点偏移量,读取下一条可用消息 final Optional> next
这篇文章对您有帮助吗?
订阅66必读
每日精选科技资讯,直达你的邮箱