- 任务调度
- 工作流自动化
- 批处理
- 后端
【免费下载链接】luigi
Luigi 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.
Luigi 是一个帮助开发者构建复杂批处理作业管道的 Python 模块,它负责依赖解析、工作流管理与可视化。本文围绕 doc/workflows.rst 展开,系统讲解构建 Luigi 工作流的两大基本构件——Task类与Target类,以及决定任务如何运行的Parameter类。读完本文,你将掌握如何用纯 Python 代码定义任务、输出与依赖关系,并理解这些抽象在源码层面的实现原理,能够直接上手编写可运行的 Luigi 管道。
一、概述:工作流的三大基石
在 Luigi 中,构建一个工作流只需要理解三个核心概念:
Task(任务):工作流中的计算单元,对应"做什么";Target(目标):任务产出的资源或状态检查点,对应"产出了什么";Parameter(参数):控制单个任务如何运行的配置输入,对应"这次怎么跑"。
其中Task和Target都是抽象基类,期望子类实现若干方法;Parameter则是贯穿二者的重要概念,它决定了同一种任务在不同参数下如何被区分与调度。三者组合起来,Luigi 就能让你在纯代码中表达任意复杂的依赖关系,而不是借助笨拙的配置 DSL——这对真实世界中"很乱"的依赖场景尤其有用。
二、Target:工作流的状态表示
2.1 概念:Target 对应什么
Target类对应以下任意一种资源形态:
- 磁盘上的一个文件;
- HDFS(Hadoop 分布式文件系统)上的一个文件;
- 某种检查点(checkpoint),例如数据库中的一条记录。
从源码看,Target是一个抽象类,其唯一必须由子类实现的方法就是exists()——它返回True当且仅当该 Target 已存在(见 luigi/target.py):
class Target(metaclass=abc.ABCMeta): @abc.abstractmethod def exists(self): """Returns True if the Target exists and False otherwise.""" pass这个"存在与否"的语义是整个工作流推进的核心:Luigi 判断一个任务是否完成,靠的就是检查它的输出 Target 是否全部exists()。
2.2 内置 Target 工具箱
在实践中,你几乎不需要自己实现 Target 子类。Luigi 自带一整套开箱即用的 Target,涵盖本地文件系统、HDFS 与各类数据库:
| Target 类 | 对应资源 | 源码位置 |
|---|---|---|
luigi.file.LocalTarget | 本地磁盘文件 | luigi/local_target.py |
luigi.contrib.hdfs.target.HdfsTarget | HDFS 文件 | luigi/contrib/hdfs/target.py |
luigi.contrib.s3.S3Target | AWS S3 对象 | luigi/contrib/s3.py |
luigi.contrib.ssh.RemoteTarget | SSH 远程文件 | luigi/contrib/ssh.py |
luigi.contrib.ftp.RemoteTarget | FTP 远程文件 | luigi/contrib/ftp.py |
luigi.contrib.mysqldb.MySqlTarget | MySQL 数据库检查点 | luigi/contrib/mysqldb.py |
luigi.contrib.redshift.RedshiftTarget | Redshift 数据库检查点 | luigi/contrib/redshift.py |
此外还有contrib目录下的更多实现(如luigi/contrib/中的 GCS、Postgres、MongoDB、Salesforce 等),可按需选用。
2.3 文件系统类 Target:原子性与 open 方法
绝大多数 Target 都是"文件系统风格"的。例如LocalTarget和HdfsTarget分别映射到本地磁盘或 HDFS 上的一个文件。除此之外,它们还封装了底层操作以保证原子性(atomic writes):数据先写入临时文件,写操作成功结束后才通过os.replace之类的原子移动落位到最终路径,避免出现半成品文件污染下游任务。
这两个类都实现了open()方法,返回一个流对象:
mode='r':以只读方式打开流;mode='w':以写入方式打开流。
以LocalTarget.open的实现为例(luigi/local_target.py):
def open(self, mode='r'): rwmode = mode.replace('b', '').replace('t', '') if rwmode == 'w': self.makedirs() return self.format.pipe_writer(atomic_file(self.path)) elif rwmode == 'r': fileobj = FileWrapper(io.BufferedReader(io.FileIO(self.path, mode))) return self.format.pipe_reader(fileobj) else: raise Exception("mode must be 'r' or 'w' (got: %s)" % mode)可见写入模式走的是atomic_file——它继承自luigi.target.AtomicLocalFile,写入时先落到path-luigi-tmp-<随机数>临时文件,close()时才调用os.replace原子地移动到最终路径(luigi/local_target.py)。这正是"文件系统操作是原子的"这句话的源码证据:如果任务运行中崩溃或忘记 close,最终路径上不会留下不完整的数据。
2.4 格式支持:Gzip 与其他 Format
Luigi 自带 Gzip 支持,用法是在构造 Target 时传入format=format.Gzip。格式层由 luigi/format.py 实现,内置了多种预置格式实例:
Text = TextFormat() UTF8 = TextFormat(encoding='utf8') Nop = NopFormat() SysNewLine = NewlineFormat() Gzip = GzipFormat() Bzip2 = Bzip2Format() MixedUnicodeBytes = MixedUnicodeBytesFormat()其中GzipFormat的读写分别通过外部进程gunzip/gzip完成(luigi/format.py),且GzipFormat还支持compression_level参数。添加对其他格式的支持也非常简单:只需实现Format接口的pipe_reader与pipe_writer两个类方法(luigi/format.py),多个格式还能通过>>运算符链式组合(ChainFormat)。
下面这张图直观地展示了任务与目标的关系:任务消费其他任务产出的 Target,同时自己也产出新的 Target。
三、Task:计算发生的地方
3.1 三个关键方法
Task类在概念上更有意思,因为计算正是在这里发生。子类可以通过实现以下三个方法来改变其行为:
| 方法 | 作用 | 源码 |
|---|---|---|
run() | 任务实际执行的计算逻辑 | luigi/task.py |
output() | 声明本任务产出的 Target;任务完成与否由输出是否存在决定 | luigi/task.py |
requires() | 声明本任务依赖的其他 Task 列表 | luigi/task.py |
三者之间的关系是:requires()声明依赖 → Luigi 调度器保证所有依赖完成后才运行本任务 →run()读取依赖的输出、执行计算 → 最终产物写入output()声明的 Target。
任务的依赖用requires()方法定义,详见 tasks 文档。下图展示了两个任务之间的依赖关系。
3.2 input():依赖到输入的桥梁
每个任务用output()声明自己的输出。此外还有一个便捷方法input():它返回每个依赖任务对应的 Target 对象,本质是requires()返回结果的"取输出"映射(luigi/task.py):
def input(self): """Returns the outputs of the Tasks returned by requires()""" return getpaths(self.requires())也就是说,self.input()拿到的是"我的依赖们产出的数据流",self.output()则是"我要写出的数据流",二者通过getpaths保持结构一致(单个 Task、list、dict 均支持,见 luigi/task.py)。
3.3 完成判定与生命周期钩子
Task.complete()的默认实现非常直白:取出output()的全部 Target,只要所有 Target 都存在,任务即视为完成(luigi/task.py):
def complete(self): outputs = flatten(self.output()) if len(outputs) == 0: warnings.warn("Task %r without outputs has no custom complete() method" % self, stacklevel=2) return False return all(map(lambda output: output.exists(), outputs))注意:没有输出(且未重写 complete)的任务会被警告并视为未完成。所以对纯计算型任务,要么声明输出,要么重写complete()。与之配套的还有:
deps():调度器内部使用,返回扁平化的requires()结果(luigi/task.py);clone():基于现有实例复制一个新实例并修改部分参数,是递归依赖场景消除样板代码的利器(luigi/task.py);on_failure()/on_success():run 抛出异常或成功完成后的钩子(luigi/task.py);ExternalTask:表示由 Luigi 之外的外部流程产出数据的任务,无run实现(luigi/task.py);WrapperTask:只包装其他任务、不产出自己的输出,完成判定为"所有依赖均完成"(luigi/task.py)。
3.4 一个完整的可运行示例
仓库中的 examples/wordcount.py 完整演示了上述所有概念:InputText继承ExternalTask声明外部数据文件,WordCount通过requires()按日期区间拉取依赖、通过input()读取依赖输出、在run()中统计词频并写入output()声明的 Target:
class InputText(luigi.ExternalTask): date = luigi.DateParameter() def output(self): return luigi.LocalTarget(self.date.strftime('/var/tmp/text/%Y-%m-%d.txt')) class WordCount(luigi.Task): date_interval = luigi.DateIntervalParameter() def requires(self): return [InputText(date) for date in self.date_interval.dates()] def output(self): return luigi.LocalTarget('/var/tmp/text-count/%s' % self.date_interval) def run(self): count = {} for f in self.input(): # input() 是 requires() 的封装,返回 Target 对象 for line in f.open('r'): for word in line.strip().split(): count[word] = count.get(word, 0) + 1 f = self.output().open('w') for word, count in count.items(): f.write("%s\t%d\n" % (word, count)) f.close() # 文件系统操作是原子的,不 close 会丢失所有数据最简任务的运行方式可参考 examples/hello_world.py 的注释,在命令行执行:
$ luigi --module examples.hello_world examples.HelloWorldTask --local-scheduler四、Parameter:任务如何被参数化
Task类对应某种"类型的作业",但在真实业务中,你通常希望对它进行参数化。例如:你的任务类每晚运行一个 Hadoop 作业生成报告,那么日期(date)显然应该是这个类的一个参数——不同日期对应不同的任务实例、不同的输出路径、不同的调度状态。
Parameter 机制的详细用法见 parameters 文档。其核心设计是:参数声明为 Task 类的类属性,Luigi 在实例化时收集所有参数并统一解析(get_params/get_param_values,见 luigi/task.py),最终为每个参数组合生成唯一的task_id,从而在调度器层面区分"同一任务类的不同参数实例"(task_id_str,见 luigi/task.py)。
Luigi 内置了丰富的参数类型,例如 examples/wordcount.py 中使用的DateParameter、DateIntervalParameter,以及数值、列表、字典、枚举等参数类型,均可直接在任务中声明使用。
五、用代码表达任意依赖
综合使用 Task、Target 与 Parameter,Luigi 让你在代码中而非某种别扭的配置 DSL 中表达任意依赖。这一点非常实用,因为真实世界的依赖关系往往很乱。以下是几类常见且难以用静态配置表达的依赖场景:
- 日期代数(date algebra):依赖关系中涉及日期推算,如"本任务依赖前一天/上一周的任务输出";
- 递归(recursion):任务依赖同类的先前实例,例如逐日回填(backfill)时
Task(date - 1)依赖链; - 枚举(enum):依赖根据枚举取值动态展开,例如对不同报表类型分别生成子任务。
(以上三图源自 2014 年底 NYC Data Science 聚会上 Erik Bernhardsson 的 Luigi 演讲。)
由于依赖是普通 Python 表达式,你可以自由组合requires()、clone()、列表推导与if/else逻辑来构造任意形状的依赖图,而这一切最终都会被 Luigi 的调度器通过deps()扁平化后解析执行。
六、总结
Luigi 工作流的全部精髓可以浓缩为三个抽象:
- Target表示状态与产物,唯一职责是实现
exists(),并借助open()与 Format 机制提供带原子性的读写能力; - Task承载计算,通过
requires()/output()/run()三驾马车驱动工作流推进,input()打通依赖输出到本任务输入的通道; - Parameter让同一任务类可被参数化复用,是构造海量相似任务实例的基础。
理解这三者以及它们组合出来的"代码即依赖"哲学,你就能利用 Luigi 构建出清晰、健壮、可回溯的批处理管道。更深入的进阶主题可继续阅读仓库内的 tasks 文档(任务编写细节)与 parameters 文档(参数定义与解析规则)。
- 任务调度
- 工作流自动化
- 批处理
- 后端
【免费下载链接】luigi
Luigi 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.
相关推荐
Luigi核心组件详解:Task、Target与参数系统
Luigi核心组件详解:Task、Target与参数系统 本文深入解析Luigi工作流框架的三大核心组件:Task类、Target系统和参数系统。Task类通过
任务调度工作流自动化批处理后端深入理解Spotify Luigi工作流构建机制
深入理解Spotify Luigi工作流构建机制 前言 在现代数据处理领域,构建可靠、可维护的工作流系统至关重要。Spotify开源的Luigi框架为解决这一需
任务调度工作流自动化批处理后端WebDataset单元测试:确保数据加载器正确性的最佳实践
WebDataset单元测试:确保数据加载器正确性的最佳实践 WebDataset是一个基于Python的高性能I/O系统,专为大型(和小型)深度学习问题设计,
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考