标记机制
标记(Marker)的用途
写入操作可能在完成之前失败,从而在存储上留下部分或损坏的数据文件。标记用于跟踪并清理所有部分完成或失败的写入操作。当写入操作开始时,会创建一个标记,表示文件写入正在进行中;当写入提交成功时,标记会被删除;如果写入操作中途失败,则会留下一个标记,表明文件不完整。使用标记的两个重要操作包括:
删除重复/部分数据文件:
- 在 Spark 中,Hudi 写客户端将数据文件的写入委托给多个执行器。某个执行器可能导致任务失败,留下已写入的部分数据文件,此时 Spark 会重试该任务,直到成功为止。
- 当启用推测执行时,同一份数据可能会被多次成功写入到不同的文件中,但最终只有一个文件会被交给 Spark driver 进程进行提交。标记有助于高效识别已写入的部分数据文件——这些文件中包含的数据与稍后成功尝试所写入的数据文件存在重复——并在提交最终确定时清理这些重复的数据文件。
回滚失败的提交:如果写入操作失败,下一个写客户端会在继续进行新写入之前回滚失败的提交。回滚过程借助标记来识别作为失败提交一部分而写入的数据文件。
如果没有标记来跟踪每次提交对应的数据文件,我们就不得不列出文件系统中的所有文件,将其与时间线中出现的文件进行关联,然后删除属于部分写入失败的那些文件。可想而知,在一个规模非常大的数据湖部署中,这样做成本会非常高。
标记结构
每个标记条目由三部分组成:数据文件名、标记扩展名(.marker)以及创建该文件的 I/O 操作类型(CREATE 表示插入,MERGE 表示更新/删除,APPEND 表示两者均可)。例如,标记 91245ce3-bb82-4f9f-969e-343364159174-0_140-579-0_20210820173605.parquet.marker.CREATE 表明对应的数据文件是 91245ce3-bb82-4f9f-969e-343364159174-0_140-579-0_20210820173605.parquet,且 I/O 类型为 CREATE。
标记写入选项
写入标记有两种方式:
- 直接将标记写入存储,这是一种遗留配置。
- 将标记写入 Timeline Server,由 Timeline Server 对标记请求进行批量处理后再写入存储(默认方式)。如下所述,这种方式可提升大文件的写入性能。
直接写入标记
直接写入存储时,会为每个数据文件创建一个对应的新标记文件,标记文件的命名如上所述。标记文件不含任何内容,即为空文件。每个标记文件都按照相同的目录层级写入存储,即提交瞬时(commit instant)和分区路径,位于 Hudi 表基础路径下的临时文件夹 .hoodie/.temp 中。例如,下图展示了向 Hudi 表写入数据时创建的标记文件及其对应数据文件的示例。在获取或删除所有标记文件路径时,该机制会先列出临时文件夹 .hoodie/.temp/<commit_instant> 下的所有路径,然后再执行相应操作。

尽管这种方式比扫描整张表以查找未提交的数据文件高效得多,但随着需要写入的数据文件数量增加,需要创建的标记文件数量也会随之增加。对于需要写入大量数据文件(例如一万个或更多)的大规模写入场景,这可能成为云存储(如 AWS S3)的性能瓶颈。在 AWS S3 中,每次创建和删除文件的调用都会触发一个 HTTP 请求,并且对存储桶中每个前缀每秒可处理的请求数量存在速率限制。当并发写入的数据文件数量和标记文件数量都非常庞大时,标记文件操作可能在写入过程中占用相当可观的时间,有时会达到数分钟甚至更长。
时间线服务器标记(默认)
为了解决上述 AWS S3 速率限制导致的性能瓶颈,我们引入了一种利用时间线服务器的新标记机制,针对文件 I/O 延迟不可忽视的存储优化了与标记相关的延迟。在下图中可以看到,基于时间线服务器的标记机制将标记创建及其他标记相关操作从各个执行器委托给时间线服务器进行集中处理。时间线服务器会批量接收标记创建请求,并以可配置的批处理间隔(默认 50 毫秒)将标记批量写入文件系统中有限数量的文件中。通过这种方式,即使数据文件数量庞大,实际的文件操作次数和标记相关的延迟也能显著降低,从而提升大规模写入的性能。

每个标记创建请求在 Javalin 时间线服务器中都是异步处理的,并在处理前先入队排队。对于每个批处理间隔,时间线服务器从队列中拉取待处理的标记创建请求,并以轮询方式将所有标记写入下一个文件。在时间线服务器内部,这种批处理是多线程的,其设计与实现旨在保证一致性与正确性。批处理间隔和批处理并发度都可以通过写入选项进行配置。

请注意,工作线程始终会通过将请求中的标记名称与时间线服务器维护的所有标记的内存副本进行比较,来检查该标记是否已经创建。存储标记的底层文件只会在第一个标记请求到来时才被读取(延迟加载)。请求的响应只有在新标记被刷新到文件之后才会返回,因此即使时间线服务器发生故障,它也能够恢复已创建的标记。这些机制确保了存储与内存副本之间的一致性,并提升了处理标记请求的性能。
注意: HDFS 尚不支持基于时间线的标记,不过由于文件系统元数据被高效地缓存在内存中,且不会遇到与 S3 相同的速率限制,用户几乎不会察觉到直接标记所带来的性能问题。
标记配置参数
| 属性名称 | 默认值 | 含义 |
|---|---|---|
hoodie.write.markers.type | timeline_server_based | 要使用的标记类型。支持两种模式:(1) direct:由执行器直接创建与每个数据文件一一对应的标记文件;(2) timeline_server_based:所有标记操作均由作为代理的时间线服务处理。新的标记条目会被批量处理,并存储在数量有限的底层文件中以提升效率。 |
hoodie.markers.timeline_server_based.batch.num_threads | 20 | 时间线服务上用于批量处理标记创建请求的线程数。 |
hoodie.markers.timeline_server_based.batch.interval_ms | 50 | 标记创建批量处理的批量间隔(毫秒)。 |
博客
评论
登录后参与评论
KnowForge