pwwang/pipen · 文件 下载 ZIP
文件最后提交记录最后更新时间
README.md
以下内容由 AI 翻译,如有问题请点此提交 issue 反馈
Python 的管道框架
Documentation | ChangeLog | Examples | API
为什么选择 pipen?
pipen 专为数据科学家、生物信息学家和研究人员设计,旨在无需传统工作流系统的复杂性,即可创建可复现、可扩展的计算流水线。
目标受众
- 数据科学家:通过自动并行化和缓存处理大型数据集
- 生物信息学家:为基因组数据构建可复现的分析流水线
- 研究人员:为计算研究创建透明、可复现的工作流
- DevOps 工程师:跨不同调度器(SLURM、SGE、Google Cloud)编排批处理任务
主要优势
1. 零配置
- 凭借合理的默认值立即开始使用
- 仅在需要时配置所需内容
- 基于配置文件的不同环境配置
2. 内置可复现性
- 基于输入/输出签名的自动作业缓存
- 流水线运行和参数的完整审计跟踪
- 依赖跟踪确保进程按正确顺序运行
3. 灵活调度
- 本地运行用于开发
- 扩展至 HPC 集群(SLURM、SGE)
- 部署到云端(Google Cloud Batch、SSH)
- 在容器中运行以确保可复现性
4. 开发者友好
- 将流水线定义为 Python 类
- 使用熟悉的 Python 语法和工具
- 可扩展的插件系统用于自定义功能
- 丰富且信息详尽的日志记录和进度跟踪
5. 数据流管理
- 流水线阶段间自动传递数据
- 支持文件、目录和内存数据
- 内置用于转换和聚合数据的操作
与替代方案的比较
| 功能 | pipen | Snakemake | Nextflow | Airflow |
|---|---|---|---|---|
| 目标受众 | 数据科学家、生物信息学家、研究人员、DevOps | 生物信息学家 | 生物信息学家 | 数据工程师 |
| 学习曲线 | 低 | 中 | 高 | 高 |
| Python 集成 | 原生 | 有限 | 有限 | 原生 |
| 调度器支持 | 6+(本地、SGE、SLURM、SSH、容器、Gbatch) | 有限 | 有限 | 基于插件 |
| 缓存 | 内置,自动 | 手动 | 手动 | 基于插件 |
| 云原生支持 | 是(Google Cloud Batch) | 部分 | 是 | 是 |
| 交互式调试 | 是 | 有限 | 否 | 否 |
| 易用性 | 将流水线定义为 Python 类,语法熟悉 | 工作流 DSL,独立的配置文件 | Python 中的 DAG 定义,复杂的 UI | |
| 零配置 | 合理的默认值,仅配置所需项 | 许多配置选项 | 需要大量配置 | 复杂的设置 |
| 良好的日志 | 丰富、信息量大、彩色编码、进度条 | 基于文本 | 基于文本 | 基础日志 |
| 高度可扩展 | 简单的插件系统,基于钩子 | 自定义规则/脚本 | 自定义算子 | 自定义算子/提供者 |
| 数据流管理 | 内置通道操作(expand_dir, collapse_files) | 手动处理 | 通道系统 | XCom 系统 |
| 可重现性 | 内置缓存,完整的审计跟踪 | 手动 | 版本化容器 | DAG 版本控制 |
| 灵活调度 | 无需代码更改即可切换调度器 | 基于配置 | 基于配置 | 基于配置 |
安装
pip install -U pipen
快速入门
example.py
from pipen import Proc, Pipen, run
class P1(Proc):
"""Sort input file"""
input = "infile"
input_data = ["/tmp/data.txt"]
output = "outfile:file:intermediate.txt"
script = "cat {{in.infile}} | sort > {{out.outfile}}"
class P2(Proc):
"""Paste line number"""
requires = P1
input = "infile:file"
output = "outfile:file:result.txt"
script = "paste <(seq 1 3) {{in.infile}} > {{out.outfile}}"
# class MyPipeline(Pipen):
# starts = P1
if __name__ == "__main__":
# MyPipeline().run()
run("MyPipeline", starts=P1)
> echo -e "3\n2\n1" > /tmp/data.txt
> python example.py
04-17 16:19:35 I core _____________________________________ __
04-17 16:19:35 I core ___ __ \___ _/__ __ \__ ____/__ | / /
04-17 16:19:35 I core __ /_/ /__ / __ /_/ /_ __/ __ |/ /
04-17 16:19:35 I core _ ____/__/ / _ ____/_ /___ _ /| /
04-17 16:19:35 I core /_/ /___/ /_/ /_____/ /_/ |_/
04-17 16:19:35 I core
04-17 16:19:35 I core version: 1.1.16
04-17 16:19:35 I core
04-17 16:19:35 I core ╔═══════════════════════════ MYPIPELINE ════════════════════════════╗
04-17 16:19:35 I core ║ My pipeline ║
04-17 16:19:35 I core ╚═══════════════════════════════════════════════════════════════════╝
04-17 16:19:35 I core plugins : verbose v1.1.1
04-17 16:19:35 I core # procs : 2
04-17 16:19:35 I core profile : default
04-17 16:19:35 I core outdir : /path/to/cwd/MyPipeline-output
04-17 16:19:35 I core cache : True
04-17 16:19:35 I core dirsig : 1
04-17 16:19:35 I core error_strategy : ignore
04-17 16:19:35 I core forks : 1
04-17 16:19:35 I core lang : bash
04-17 16:19:35 I core loglevel : info
04-17 16:19:35 I core num_retries : 3
04-17 16:19:35 I core scheduler : local
04-17 16:19:35 I core submission_batch: 8
04-17 16:19:35 I core template : liquid
04-17 16:19:35 I core workdir : /path/to/cwd/.pipen/MyPipeline
04-17 16:19:35 I core plugin_opts :
04-17 16:19:35 I core template_opts : filters={'realpath': <function realpath at 0x7fc3eba12...
04-17 16:19:35 I core : globals={'realpath': <function realpath at 0x7fc3eba12...
04-17 16:19:35 I core Initializing plugins ...
04-17 16:19:36 I core
04-17 16:19:36 I core ╭─────────────────────────────── P1 ────────────────────────────────╮
04-17 16:19:36 I core │ Sort input file │
04-17 16:19:36 I core ╰───────────────────────────────────────────────────────────────────╯
04-17 16:19:36 I core P1: Workdir: '/path/to/cwd/.pipen/MyPipeline/P1'
04-17 16:19:36 I core P1: <<< [START]
04-17 16:19:36 I core P1: >>> ['P2']
04-17 16:19:36 I verbose P1: in.infile: /tmp/data.txt
04-17 16:19:36 I verbose P1: out.outfile: /path/to/cwd/.pipen/MyPipeline/P1/0/output/intermediate.txt
04-17 16:19:38 I verbose P1: Time elapsed: 00:00:02.051s
04-17 16:19:38 I core
04-17 16:19:38 I core ╭═══════════════════════════════ P2 ════════════════════════════════╮
04-17 16:19:38 I core ║ Paste line number ║
04-17 16:19:38 I core ╰═══════════════════════════════════════════════════════════════════╯
04-17 16:19:38 I core P2: Workdir: '/path/to/cwd/.pipen/MyPipeline/P2'
04-17 16:19:38 I core P2: <<< ['P1']
04-17 16:19:38 I core P2: >>> [END]
04-17 16:19:38 I verbose P2: in.infile: /path/to/cwd/.pipen/MyPipeline/P1/0/output/intermediate.txt
04-17 16:19:38 I verbose P2: out.outfile: /path/to/cwd/MyPipeline-output/P2/result.txt
04-17 16:19:41 I verbose P2: Time elapsed: 00:00:02.051s
04-17 16:19:41 I core
MYPIPELINE: 100%|██████████████████████████████| 2/2 [00:06<00:00, 0.35 procs/s]
> cat ./MyPipeline-output/P2/result.txt
1 1
2 2
3 3
示例
更多示例见 examples/,以及一个更贴近实际案例的示例见:
https://github.com/pwwang/pipen-report/tree/master/example
插件库
插件让 pipen 更加出色。
pipen-annotate: 使用 docstring 为 pipen 进程添加注释pipen-args: pipen 的命令行参数解析器pipen-board: 在 Web 上可视化 pipen 流水线的配置和运行pipen-diagram: 为 pipen 绘制流水线图pipen-dry: pipen 流水线的试运行器pipen-filters: 为 pipen 模板添加一组有用的过滤器。pipen-lock: pipen 的进程锁,防止同时多次运行。pipen-log2file: 将 pipen 的运行日志保存到文件pipen-poplog: 将作业日志填充到流水线的运行日志中pipen-report: 为 pipen 生成报告pipen-runinfo: 将 pipen 的运行信息保存到文件pipen-verbose: 在 pipen 的日志中添加详细信息。pipen-email: 发送流水线状态变更的电子邮件通知。pipen-gcs: 一个用于处理 Google Cloud Storage 中文件的 pipen 插件。pipen-deprecated: 一个用于将进程标记为已弃用的 pipen 插件。pipen-mcp: 一个将 pipen 进程转换为 MCP (model context protocol) 进程的 pipen 插件。pipen-cli-init: 一个用于创建 pipen 项目(流水线)的 pipen CLI 插件pipen-cli-ref: 为进程创建参考文档pipen-cli-require: 一个用于检查流水线要求的 pipen cli 插件pipen-cli-run: 一个用于运行进程或流水线的 pipen cli 插件pipen-cli-gbatch: 一个用于将流水线提交到 Google Batch Jobs 的 pipen cli 插件