如何定义你的第一个Luigi任务:requires、output与run方法完全指南

发布时间:2026/9/18 10:13:01
如何定义你的第一个Luigi任务:requires、output与run方法完全指南 如何定义你的第一个Luigi任务requires、output与run方法完全指南【免费下载链接】luigiLuigi is a Python module that helps you build complex pipelines of batch jobs. It handles dependency resolution, workflow management, visualization etc. It also comes with Hadoop support built in.项目地址: https://gitcode.com/gh_mirrors/lu/luigiLuigi 是一个 Python 模块帮助你构建复杂的批处理任务管道batch pipeline。它自动处理依赖解析、任务调度、可视化与失败重试还内置了 Hadoop 支持。本文面向新手带你从零定义第一个 Luigi 任务彻底弄懂requires、output、run这三个核心方法的作用与协作方式。一个Luigi任务的三大构件每个 Luigi 任务都是luigi.Task的子类通常包含三个部分示意图见上方构件作用一句话理解requires(self)声明依赖的上游任务我要先等谁完成output(self)返回一个 Target 对象我的产物是什么run(self)真正的业务逻辑我怎么生产产物核心思想是用输出产物代替任务本身来管理依赖。Luigi 通过检查 output 是否已存在来判断任务是否完成存在就跳过这就是所谓幂等idempotent机制。requires 方法声明任务依赖requires返回一个或多个上游 Task 对象Luigi 会自动把它们排进执行顺序。它可以返回单个任务、列表、字典等任意结构来源定义见 luigi/task.pydef requires(self): return GenerateWords()如果依赖的是外部系统产出的文件别人写的文件不是 Luigi 任务写的建议用luigi.ExternalTask包一层只实现output即可官方示例见 examples/wordcount.pyclass InputText(luigi.ExternalTask): date luigi.DateParameter() def output(self): return luigi.LocalTarget(/var/tmp/text/%Y-%m-%d.txt % date)⚠️ 注意requires里不能直接返回 Target 对象必须包成 Task。output 方法定义任务产物output返回一个或多个luigi.target.Target对象最常用的是本地文件luigi.LocalTarget定义在 luigi/local_target.py也可以换成 HDFS、S3 等远程目标def output(self): return luigi.LocalTarget(/tmp/result.txt)官方推荐一个任务只返回一个 Target因为多个输出会破坏原子性。run 方法编写业务逻辑run里放真正的代码。它通过两个便捷方法拿数据self.input()相当于把 requires 里的任务换成它们的 output返回 Target 对象self.output()拿到自己的产物目标用.open(w)写入。一个最小可运行的字统计任务import luigi class CountWords(luigi.Task): def requires(self): return GenerateWords() def output(self): return luigi.LocalTarget(/tmp/word_count.txt) def run(self): with self.input().open(r) as f: words f.read().splitlines() with self.output().open(w) as f: f.write(total: %d\n % len(words))关键约定run执行完毕后complete()必须返回 True——也就是说 run 结束时产物必须已落盘否则 Luigi 会报 Unfulfilled dependencies at run time。完整示例从输入到可视化的词频管道examples/wordcount.py 展示了经典三步结构外部文本任务InputText只有 output→ 词频任务WordCountrequires output run 俱全。WordCount还支持日期区间参数一次可以扇出成多个子任务并行跑。运行后你可以启动内置的 Web 可视化工具直观看到整条管道的依赖与执行状态红失败、蓝运行中、黄等待、绿完成命令行运行你的任务写好任务后用luigi命令行工具启动详见 doc/running_luigi.rstluigi --module my_module MyTask --x 123 --local-scheduler参数名中的下划线要写成连字符如--my-parameter且参数必须放在任务名之后。常见新手问题速查Q1任务重复执行怎么办先检查 output 路径是否写对只要产物存在Luigi 就不会重跑该任务。Q2依赖在运行时才能确定在run中yield OtherTask()实现动态依赖但注意run会被重新从头执行务必保持幂等。参考 examples/dynamic_requirements.py。Q3写二进制文件报编码错给LocalTarget传formatNop避免默认的文本编解码干扰。Q4多输入任务如何取数self.input()会保持与requires相同的结构列表、字典照原样按需下标访问即可。小结requires管依赖谁output管产出什么run管怎么干依赖管理以output 产物为中心天然幂等、可断点续跑上手顺序读 doc/tasks.rst → 跑通 examples/hello_world.py → 复现 examples/wordcount.py → 用可视化工具验证依赖图。掌握这三件套你就迈出了构建生产级数据管道的第一步 【免费下载链接】luigiLuigi is a Python module that helps you build complex pipelines of batch jobs. It handles dependency resolution, workflow management, visualization etc. It also comes with Hadoop support built in.项目地址: https://gitcode.com/gh_mirrors/lu/luigi创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考