Luigi 工作流构建指南:深入理解 Task、Target 与 Parameter 三大核心抽象
2026/9/20 16:30:23 网站建设 项目流程
  • 任务调度
  • 工作流自动化
  • 批处理
  • 后端

【免费下载链接】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.

项目地址:https://gitcode.com/gh_mirrors/lu/luigi
点击查看免费下载

Luigi 是一个帮助开发者构建复杂批处理作业管道的 Python 模块,它负责依赖解析、工作流管理与可视化。本文围绕 doc/workflows.rst 展开,系统讲解构建 Luigi 工作流的两大基本构件——Task类与Target类,以及决定任务如何运行的Parameter类。读完本文,你将掌握如何用纯 Python 代码定义任务、输出与依赖关系,并理解这些抽象在源码层面的实现原理,能够直接上手编写可运行的 Luigi 管道。

一、概述:工作流的三大基石

在 Luigi 中,构建一个工作流只需要理解三个核心概念:

  • Task(任务):工作流中的计算单元,对应"做什么";
  • Target(目标):任务产出的资源或状态检查点,对应"产出了什么";
  • Parameter(参数):控制单个任务如何运行的配置输入,对应"这次怎么跑"。

其中TaskTarget都是抽象基类,期望子类实现若干方法;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.HdfsTargetHDFS 文件luigi/contrib/hdfs/target.py
luigi.contrib.s3.S3TargetAWS S3 对象luigi/contrib/s3.py
luigi.contrib.ssh.RemoteTargetSSH 远程文件luigi/contrib/ssh.py
luigi.contrib.ftp.RemoteTargetFTP 远程文件luigi/contrib/ftp.py
luigi.contrib.mysqldb.MySqlTargetMySQL 数据库检查点luigi/contrib/mysqldb.py
luigi.contrib.redshift.RedshiftTargetRedshift 数据库检查点luigi/contrib/redshift.py

此外还有contrib目录下的更多实现(如luigi/contrib/中的 GCS、Postgres、MongoDB、Salesforce 等),可按需选用。

2.3 文件系统类 Target:原子性与 open 方法

绝大多数 Target 都是"文件系统风格"的。例如LocalTargetHdfsTarget分别映射到本地磁盘或 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_readerpipe_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 中使用的DateParameterDateIntervalParameter,以及数值、列表、字典、枚举等参数类型,均可直接在任务中声明使用。

五、用代码表达任意依赖

综合使用 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 工作流的全部精髓可以浓缩为三个抽象:

  1. Target表示状态与产物,唯一职责是实现exists(),并借助open()与 Format 机制提供带原子性的读写能力;
  2. Task承载计算,通过requires()/output()/run()三驾马车驱动工作流推进,input()打通依赖输出到本任务输入的通道;
  3. 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.

项目地址:https://gitcode.com/gh_mirrors/lu/luigi
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询