← 文章 / 数据与数据库
InfoQ 2小时前 · 2026-09-03 12:33:14 · 1 阅读

超越偏移量延迟:PB 级规模下 Hudi 数据湖流水线的队列等待时间计算

问题背景

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<TopicPartition, Long> checkpointOffsets = hudiResult.getPartitionToCheckpoint()    .entrySet().stream()    .collect(Collectors.toMap(        e -> new TopicPartition(topic, e.getKey()),        Map.Entry::getValue    ));// 步骤2、3:定位到检查点偏移量,读取下一条可用消息final Optional<ConsumerRecord<String, String>> nextRecord = kafkaClient.nextRecord(topic, checkpointOffsets);// 步骤4:计算队列等待时间(time‑in‑queue)nextRecord.ifPresent(record -> {    long currentTimeMs = OffsetDateTime.now(Clock.systemUTC()).toInstant().toEpochMilli();    Duration lag = Duration.ofMillis(Math.max(0L, currentTimeMs - record.timestamp()));    long lagSeconds = lag.getSeconds();    // 将lagSeconds上报到指标监控系统});
复制代码

从 S3 获取 Hudi 检查点

关键在于遍历 Hudi 活跃时间线,找到包含检查点元数据的最新一次提交。我们使用 Apache Hudi SDK 的 HoodieTableMetaClient 从 S3 读取 .hoodie/ 时间线目录。

HoodieTableMetaClient 底层使用 Hadoop 的 S3A 文件系统访问 S3。它需要一份 Hadoop 配置,我们从报告器运行的 Spark 会话中获取配置。这里不涉及实际的 Spark 数据处理;Spark 在这里纯粹作为运行时环境,为 Hudi SDK 提供其 S3 访问层。

public static HoodieResult findLatestCommitWithCheckpoint(String datasetName, String basePath, int maxCommits) {    // getActiveTimeline 会基于表的S3路径初始化HoodieTableMetaClient,并返回其活跃时间线    HoodieActiveTimeline timeline = getActiveTimeline(basePath);    HoodieTimeline commits = timeline.getCommitsTimeline().filterCompletedInstants();    // maxCommits 从环境变量 MAX_COMMIT_DEPTH 读取,默认值为100    int depth = Math.min(maxCommits, commits.countInstants());        // n=0代表最新的提交,按时间倒序遍历    for (int n = 0; n < depth; n++) {           HoodieInstant commit = commits.nthFromLastInstant(n).get();        HoodieCommitMetadata metadata = HoodieCommitMetadata.fromBytes(            timeline.getInstantDetails(commit).get(), HoodieCommitMetadata.class);          // CHECKPOINT_KEY = "deltastreamer.checkpoint.key"        if (metadata.getMetadata(CHECKPOINT_KEY) != null) {            return new HoodieResult(metadata, commit.getTimestamp(), datasetName, n);        }    }    return new HoodieResult(-1); // 未找到检查点}
复制代码

检查点字符串的格式为:topicName,0:offset0,1:offset1,...,每个分区一个条目。关键点在于,Hudi 保存的是下一次要读取的偏移量,即下一次运行的恢复点,而不是最后消费的偏移量。所以如果 Hudi 最后提交了偏移量 1199 的消息,它保存的是 1200。Kafka 消费者定位到偏移量 1200 时,就处于数据湖中尚未提交的第一条消息的位置。

在 Kafka 中定位到检查点

分配分区,将每个分区定位到各自的检查点偏移量,并返回所有分区中最早的记录。

// 分配分区并定位到检查点偏移量consumer.assign(checkpointOffsets.keySet());checkpointOffsets.forEach(consumer::seek);// 拉取消息,找出所有分区中时间最早的记录ConsumerRecord<String, String> earliest = null;ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));for (ConsumerRecord<String, String> record : records) {    if (earliest == null || record.timestamp() < earliest.timestamp()) {        earliest = record;    }}// earliest.timestamp() 即为 X — 第一条未消费消息的时间戳
复制代码

上面介绍的三个代码片段就是整个算法的全部内容。跑通正常路径很简单;要让这个指标在生产环境中稳定可靠,关键在于处理一系列边界情况,下面逐一展开说明。

边界情况与注意事项

跨生产者的时钟偏移

由于延迟指标使用的是生产者设置的消息时间戳,如果某些生产者的系统时钟向前漂移,就会使延迟看起来异常偏低(甚至为负)。算法使用 Math.max(0L, currentTimeMs - record.timestamp()) 将延迟下限设为 0。在取下限之前,如果原始计算值持续为负,这本身就是值得单独告警的信号——它意味着存在需要排查的时钟偏移问题。

多分区主题

在大规模场景下,Kafka 主题有很多分区。Hudi 中保存的检查点包含每个分区的一个偏移量(如上面的示例所示)。算法将每个分区定位到各自的检查点偏移量,并拉取下一条可用记录。每个分区提供一个候选消息,然后我们取所有分区中时间戳最早的那条记录,而不是平均值,也不是最新值。这条消息代表了仍在等待被提交到数据湖的最旧数据,这是最坏情况下的延迟,也是 SLA 需要真正关注的数字。取平均值会掩盖某个分区处理缓慢的问题;取最新值则会掩盖某个分区停滞不前的问题。

多 Hudi 表写入端与 Kafka 检查点

在生产环境中,一个 Hudi 表可能有多个写入端,而并非所有写入端都会提交 Kafka 检查点元数据。这种情况常见于迁移期间,当旧管道和新管道同时向同一张表写入,而旧管道正处于逐步下线的过程中。

我们遇到的正是这种情况。旧架构里一共有两条管道:一条从 Kafka 读取数据,把原始数据写入 S3;另一条读取 S3 中的输出,将数据转换后写入 Hudi 表。

第二个管道的数据源是 S3,而不是 Kafka。因此它的提交消息不会有 deltastreamer.checkpoint.key。我们的新框架直接从 Kafka 读取数据并写入同一张 Hudi 表,它的提交包含了检查点键。

在迁移重叠期间,两条管道同时向同一张表写入。当最近一次提交来自旧管道时,原算法找不到任何检查点元数据——没有可供定位的偏移量,于是延迟只能默认取七天的上限值,在延迟图表中产生了虚假的峰值,而这些并非真实的数据新鲜度故障。

图 3. 共享 Hudi 表的多个写入端。

修复方案

根本原因在于,原算法只查看最近一次提交。真正的解决办法是:不再默认最近的提交就是正确的那个。我们改进了算法,让它按时间倒序沿时间线往回遍历,最多回退 MAX_COMMIT_DEPTH 次提交,直到找到真正包含 deltastreamer.checkpoint.key 的最近一次提交。由于只有我们的新 Kafka 数据源管道会写入这个键,往回遍历时会自动跳过旧管道的提交,最终落在有效的检查点上。这个方案无需引入任何外部状态,就消除了虚假尖峰。

如果在搜索深度内找不到任何包含检查点的提交,报告器会直接抑制该指标的上报,而不是发布一个可能产生误导的值。指标缺失会在告警系统中呈现为“No Data”告警——相比一个看似真实、实则虚假的数字,这才是更诚实的信号。

图 4. 逆序遍历 Hudi 时间线,找到有效的检查点。

更新后的算法

从 Hudi 时间线获取最新提交,并按逆时间顺序回退遍历。找到包含 deltastreamer.checkpoint.key 的最新提交;如果在 MAX_COMMIT_DEPTH 内未找到,则完全抑制指标的上报。

  • 在每个分区中定位到该偏移量

  • 读取所有分区中最早的消息并获取时间戳 X

  • 延迟 = 当前时间戳 - X

  • 上限设为七天

真实的结果

指标报告器通过 EventBridge 在 EMR Serverless 上运行,每十五分钟运行一次。每个管道通过 YAML 配置中的 slaInMinutes 字段定义自己的 SLA 阈值:

metadata:  # 当数据过期超过30分钟时触发告警  slaInMinutes: 30
复制代码

如果省略了该字段,流式管道的默认值为 60 分钟,批处理管道为 1440 分钟(24 小时)。报告器在内部将 slaInMinutes 转换为秒。SLA 违反情况以 lagSeconds / slaThresholdSeconds 的比值上报,即一个介于 0.0 和 1.0 之间的比率(上限为 1.0),其中 1.0 表示完全违反。处于 0.7 的管道表示已达到其阈值的百分之七十。仪表盘展示的是向违反阈值逼近的过程,而不是从绿色到红色的二值跳变。现有管道无需做任何更改;报告器纯粹是一个观察层。

上述算法是生产环境多轮迭代打磨的成果。首次部署时,我们立刻撞上了一堆未曾预料到的边界情况。每一轮迭代都伴随着误报或指标静默失效的问题。以下是我们从中总结的经验。

纪元时间戳陷阱

这个问题追踪起来非常头疼。我们的第一版计算延迟的方式是 latestKafkaRecord.timestamp() - hudiCommitTimestamp,假设随着 Hudi 滞后,这个差距只会扩大。我们没有意识到,当 Hudi 无法找到包含检查点元数据的提交时,它不会报错。相反,它会静默地默认使用纪元零的时间戳 19700101000000000。这个 yyyyMMddHHmmssSSS 格式的字符串转换为毫秒时值为零,导致我们的指标爆炸成一个天文数字,毫无意义。我们设置的 Math.max(0L,) 下限在这里毫无用处。这次换来的深刻教训是,我们必须显式地防范纪元默认值,并调整我们的逻辑。改为 currentTime - timestamp_of_first_unconsumed_message (当前时间减去第一条未消费消息的时间戳)之后,我们终于用对了指标,这证明了有时候缺失的指标比虚构的数值更好。

最新提交并不总是正确的提交

原始算法只是读取 Hudi 时间线中的最新提交。在迁移重叠期间,当旧管道和新管道同时向同一张表写入时,最新提交经常是旧管道的,而它没有检查点键。改为基于深度的回退查找(最多回退 MAX_COMMIT_DEPTH 个提交,默认为 100)解决了这个问题。新算法会跳过没有检查点键的提交,找到最新的有效提交。我们还将提交深度(即我们不得不回退查找了多少个提交)作为一个单独的指标进行报告。提交深度偏高是一个领先指标——即使延迟指标还没越过阈值,它也已经能提前预示主管道出问题了。

缺失的 Kafka 时间戳

Kafka 允许生产者省略消息时间戳。当这种情况发生时,记录不会携带空值。Kafka 使用哨兵值 -1()来表示缺失的时间戳。当算法遇到时间戳为 -1 的记录时,它会尝试用它来计算延迟并发布一个垃圾指标。修复方案是显式检测 -1 时间戳并跳过该管道的报告。与纪元时间戳的情况一样,沉默比会导致误导的数字更有用。

SLA 作为比率,而非二元值

“SLA 比率”并不是一个行业术语,而是我们在生产环境中摸索出来的一个设计选择。我们最初的 SLA 指标是个二元值:要么违反,要么没违反。问题在于,它只会在 SLA 已经被打破之后才触发。工程师们想知道的是管道何时已经用掉了阈值的百分之七八十,这样他们就能在客户受到影响之前介入排查。改用 0.0–1.0 的比率后,团队就能得到一个领先信号,可以把 0.7 设为警告、1.0 设为严重告警,而不是一个单一的全有或全无的阈值。这种做法与 SRE 跟踪错误预算消耗速率的做法如出一辙——你要观察的是预算消耗的速度,而不是等到预算耗尽之后才做出反应。

未来展望

如果你想套用这种模式,需要具备三样东西:

  • 能访问 Hudi 表在 S3 上的 .hoodie/ 目录;

  • 一个可以定位任意偏移量的 Kafka 客户端;

  • 一个定期运行报告器的调度器。

核心算法适用于任何写入 S3 的 Hudi Delta Streamer 管道。如果你使用的是 Delta Lake + 结构化流式处理,Spark 会将 Kafka 偏移量存储在流式检查点目录中,而不是表元数据本身中,但原理相同:从框架持久化偏移量的任何地方找到最后提交的 Kafka 偏移量,定位到该位置,计算时间戳差值。

在我们的路线图中,我们正在评估将这种方法应用于 Iceberg 管道,因为我们正在从 Hudi 迁移出来。我们还在探索将此指标与异常检测库配对,例如开源的 ProphetLuminaire,而不是仅依赖固定的 SLA 阈值。

常见问题

指标上报器会干扰 HoodieStreamer 的消费者群组吗?不会。上报器使用一个专门的消费组 ID,与 HoodieStreamer 自身的消费者群组完全隔离。它还设置了 enable.auto.commit=false,因此永远不会把偏移量提交回 Kafka。它只是读取数据然后丢弃。不存在引发再均衡或干扰摄入管道的风险。

如果某个主题分区是空的会怎样?如果在定位到检查点偏移量之后,consumer.poll() 没有返回任何记录,上报器会将其视为空主题,并上报零延迟。它不会挂起。五百毫秒的轮询超时限制了等待时间,空结果也会被显式处理。延迟默认取零,并且仍然会正常上报到指标系统。

查看英文原文:https://www.infoq.com/articles/beyond-offset-lag-kafka-apache-hudi/

原始来源: InfoQ

评论 (0)