Tar_包递归解压
更新时间:2026-07-29
简介
Tar 包递归解压算子(RecursiveTarUncompress),基于 stdlib tarfile 递归解压可能嵌套的 tar 归档,是异构数据入湖链路的第一步。带 tar-slip(路径穿越)防护,无第三方依赖,离线 CPU 运行。
功能描述
• 将顶层 tar 归档解压到指定输出目录 output_path;若归档内还含 tar,则在原地继续解压,最多递归 max_depth 层。
• 嵌套 tar 解压完成后原地删除并去重,保证内容都落在输出目录根下,避免多余子目录。
• 内置路径穿越防护:解压前校验每个成员路径必须位于目标目录内,越界成员直接拦截报错。
• 逐行处理并做异常隔离:成功返回 "SUCCESS:<output_path>",失败返回 "Failed: <错误信息>";内部维护提交/成功/失败计数。
算子参数
输入
| 输入 | 含义 |
|---|---|
| input_col | 待解压的 tar 归档文件路径(.tar / .tar.gz / .tgz) |
| output_col | 解压目标目录路径(不存在时自动创建) |
输出
| 输出 | 含义 |
|---|---|
| result | large_string 类型;成功为 "SUCCESS:<output_path>",失败为 "Failed: <错误信息>" |
参数
| 参数名称 | 类型 | 默认值 | 描述 |
|---|---|---|---|
| max_depth | int | 3 |
递归解压的最大层数,用于限制嵌套 tar 的展开深度(会被转为 int) |
调用示例
Python
1from __future__ import annotations
2
3import os
4
5import daft
6from daft import col
7
8from daft.aihc.common.udf import aihc_udf
9from daft.aihc.functions.process.recursive_tar_extractor_udf import RecursiveTarUncompress
10
11INPUT_PATH = "/mnt/pfs/xxx/outer.tar" # MOCK
12OUTPUT_PATH = "/mnt/pfs/xxx/out" # MOCK
13
14if __name__ == "__main__":
15 if os.getenv("DAFT_RUNNER", "native") == "ray":
16 import ray
17
18 ray.init(ignore_reinit_error=True)
19 daft.set_runner_ray()
20 daft.set_execution_config(min_cpu_per_task=0)
21
22 ds = daft.from_pydict({"src": [INPUT_PATH], "dst": [OUTPUT_PATH]})
23 ds = ds.with_column(
24 "result",
25 aihc_udf(
26 RecursiveTarUncompress,
27 construct_args={"max_depth": 3},
28 num_cpus=1,
29 concurrency=1,
30 batch_size=1,
31 )(col("src"), col("dst")),
32 )
33 ds.show()
评价此篇文章
