第一章:Prefect 概述与核心理念
了解 Prefect 是什么、能做什么、为什么选择它,建立对工作流自动化工具的整体认知。
1.1 什么是 Prefect?
| 名称 | 说明 | 注意事项 |
|---|
| Prefect | 一个现代的、开源的 Python 工作流(数据流)自动化框架,用于构建、调度和监控数据管道。它强调开发者体验,支持动态工作流、丰富的状态追踪和灵活的部署方式。 | Prefect 不仅适用于 ETL/ELT,也适用于机器学习流水线、定时任务、批处理等场景。 |
| 开发者优先设计 | Prefect 以 Python 原生代码为核心,允许用户使用标准 Python 语法定义任务和流程,无需 DSL(领域特定语言)。 | 降低了学习门槛,提升可维护性。 |
| 动态执行引擎 | 支持在运行时动态生成任务(如 map 操作),不同于 Airflow 的静态 DAG 构建。 | 适合处理不确定数量输入的任务(如文件列表处理)。 |
| 开源 + 托管服务 | 提供开源版本(Prefect OSS)和云服务(Prefect Cloud),支持自托管(Prefect Server)。 | 可根据团队规模选择部署模式。 |
1.2 Prefect 的核心优势与适用场景
| 优势/场景 | 说明 | 注意事项 |
|---|
| 简单易用 | 使用装饰器 @flow 和 @task 即可将函数转化为工作流单元,无需复杂配置。 | 特别适合数据科学家和工程师快速构建自动化流程。 |
| 强大的可观测性 | 提供详细的日志、状态追踪、UI 界面(Orion),便于调试和监控。 | 推荐配合 Prefect Cloud 或本地 Orion UI 使用。 |
| 灵活的执行模型 | 支持同步、异步、并行、分布式任务执行,可自定义 Task Runner。 | 可根据资源情况选择合适的执行策略。 |
| 内置重试与错误处理 | 支持任务级重试、跳过失败、条件分支等容错机制。 | 提高工作流的健壮性。 |
| 适用场景:ETL/ELT | 自动化从数据抽取、转换到加载的全过程。 | 可结合 Pandas、Dask、Spark 等工具。 |
| 适用场景:ML 工作流 | 训练、评估、部署模型的自动化流水线。 | 支持参数化训练任务。 |
| 适用场景:定时任务替代 cron | 更可靠的定时调度,支持依赖、重试、通知。 | 比 shell 脚本更易维护。 |
| 适用场景:API 编排 | 调用多个外部 API 并聚合结果。 | 支持异步请求提升效率。 |
1.3 Prefect 1.x vs 2.x 架构对比
| 对比项 | Prefect 1.x | Prefect 2.x | 注意事项 |
|---|
| 架构模式 | Client-Server 模型,依赖复杂的 Prefect Core 引擎和 Agents。 | 基于 Orion 引擎的轻量级架构,去除了 Agents,使用 Deployments 和 Work Pools。 | 2.x 更简化,部署更容易。 |
| 核心组件 | Flow, Task, Project, Agent, Executor | Flow, Task, Deployment, Work Pool, Worker | 2.x 引入 Deployment 作为部署单元。 |
| 定义方式 | 使用 Flow() 类或装饰器,语法较复杂。 | 全面使用装饰器 @flow 和 @task,纯函数式风格。 | 2.x 更 Pythonic,易于理解。 |
| 调度机制 | 使用本地或远程 Agent 拉取任务。 | 使用 Work Pool 和 Worker 拉取 Deployment 任务。 | 2.x 支持 Kubernetes、Docker、Process 等多种 Worker。 |
| 状态存储 | 默认使用本地 SQLite,可配置 Postgres。 | 同样支持 SQLite 和 Postgres,通过 Orion API 统一管理。 | 生产环境建议使用 Postgres。 |
| CLI 工具 | prefect agent, prefect auth 等命令分散。 | 统一 CLI prefect,命令结构更清晰(如 prefect deploy, prefect worker)。 | 2.x CLI 更现代化。 |
| API 设计 | GraphQL API,较难直接调用。 | RESTful API,更易集成和扩展。 | 有利于自定义监控和集成。 |
| 社区与维护 | 已停止新功能开发,仅维护。 | 当前主推版本,持续更新。 | 新项目应使用 Prefect 2.x。 |
1.4 Prefect 核心概念概览(Flow, Task, State, Result 等)
| 概念 | 说明 | 注意事项 |
|---|
| Flow | 一个由多个任务组成的有向无环图(DAG),是工作流的顶层容器。用 @flow 装饰器定义。 | 每个 Flow 是一个可调用对象,可参数化。 |
| Task | 工作流中的基本执行单元,表示一个具体操作。用 @task 装饰器定义。 | Task 可重用、可缓存、可重试。 |
| State | 表示任务或流程的执行状态,如 RUNNING, COMPLETED, FAILED, SKIPPED 等。 | 所有执行结果都封装在 State 对象中。 |
| Result | 任务执行后的返回值,可持久化存储(如本地、S3)。 | 启用 result storage 可实现跨运行共享数据。 |
| Task Runner | 控制任务如何执行的组件,如串行、并发、多进程等。 | 默认为 ConcurrentTaskRunner。 |
| Deployment | 一个可调度的 Flow 配置,包含入口点、参数、调度计划等元信息。 | 是部署到生产环境的基本单位。 |
| Work Pool | 一组 Worker 的逻辑分组,Deployment 关联到 Work Pool。 | 用于隔离不同环境或资源类型。 |
| Worker | 从 Work Pool 拉取任务并执行的长期运行进程。 | 支持 Process、Kubernetes、Docker 等类型。 |
| Block | 可复用的配置单元,用于存储连接信息(如数据库、云存储、通知服务)。 | 实现配置与代码分离,支持加密(SecretBlock)。 |
| Orion UI | Prefect 2.x 的可视化控制台,用于监控 Flow Runs、查看日志、管理部署。 | 默认运行在 http://127.0.0.1:4200。 |
第二章:环境准备与快速上手
完成本地开发环境搭建,运行第一个 Flow,体验基本流程。
2.1 安装 Prefect 2.x
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| pip 安装 | pip install prefect | 安装最新稳定版 Prefect 2.x | pip install prefect | 建议在虚拟环境中安装。 |
| 安装特定版本 | pip install prefect==2.14.0 | 安装指定版本 | pip install prefect==2.10.0 | 用于版本锁定或兼容性测试。 |
| 安装额外依赖 | pip install "prefect[extra]" | 安装带特定功能的 Prefect(如 cloud, aws, gcp) | pip install "prefect[aws]"
pip install "prefect[dashboard]" | dashboard 启用本地 Orion UI。 |
2.2 验证安装与版本检查
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 检查版本 | prefect version | 查看当前安装的 Prefect 版本 | prefect version 输出示例:Version: 2.14.0 | 确保为 2.x 系列。 |
| Python 内检查 | import prefect
print(prefect.__version__) | 在脚本中获取版本号 | import prefect
print(prefect.__version__) | 用于脚本兼容性判断。 |
| 验证 CLI 可用 | prefect --help | 查看所有可用命令 | prefect --help | 初次安装后建议执行。 |
2.3 编写你的第一个 Flow
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 导入模块 | from prefect import flow, task | 引入核心装饰器 | from prefect import flow, task | 必须先导入才能使用。 |
| 定义 Task | @task
def my_task():
return "Hello" | 将函数标记为 Prefect 任务 | @task
def say_hello():
print("Hello from task!")
return "Hi" | Task 函数可记录日志、返回值。 |
| 定义 Flow | @flow
def my_flow():
result = my_task() | 将函数标记为 Flow,组织任务调用 | @flow
def hello_flow():
result = say_hello()
print(f"Got: {result}") | Flow 是任务的容器。 |
2.4 运行 Flow 并查看结果
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 直接调用 Flow | my_flow() | 在 Python 中运行 Flow | if __name__ == "__main__":
hello_flow() | 本地调试最简单方式。 |
| 获取返回值 | result = my_flow() | 捕获 Flow 执行结果 | final_result = hello_flow()
print(final_result) | Flow 可返回数据用于后续处理。 |
| 查看日志输出 | 内置日志捕获 | Prefect 自动捕获 print 和 logging | import logging
logging.info("Task started") | 日志可在 Orion UI 查看。 |
2.5 使用 prefect CLI 基础命令
| 命令 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 启动本地服务器 | prefect orion start | 启动 Orion API 和 UI | prefect orion start | 默认监听 4200 端口。 |
| 查看帮助 | prefect --help 或 prefect <command> --help | 获取命令使用说明 | prefect flow --help | 推荐常用命令后加 --help。 |
| 登录 Prefect Cloud | prefect cloud login -k <your-key> -w <workspace> | 登录云端工作区 | prefect cloud login -k xyz123 -w my-company/ws | 仅使用 Cloud 时需要。 |
| 设置配置项 | prefect config set PREFECT_API_URL=... | 配置 API 地址等参数 | prefect config set PREFECT_LOGGING_LEVEL=DEBUG | 配置保存在 ~/.prefect/config.toml。 |
| 查看运行状态 | prefect ls | 列出当前 Flow 和 Deployment | prefect ls | 需已注册 Deployment 才可见。 |
第三章:Flow 与 Task 基础构建
掌握定义和组织任务的基本方式,理解函数如何转化为可调度单元。
3.1 使用 @flow 装饰器定义 Flow
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
@flow 基础用法 | @flow
def my_flow():
... | 将普通函数转换为 Prefect Flow,作为任务的容器 | from prefect import flow
@flow
def greet_flow():
print("Starting flow") | 必须导入 flow,函数可包含任务调用或逻辑。 |
| 命名 Flow | @flow(name="Custom Name") | 自定义 Flow 在 UI 中显示的名称 | @flow(name="Data Ingestion Pipeline")
def ingest():
... | 默认使用函数名,建议命名清晰。 |
| 日志级别设置 | @flow(log_prints=True) | 启用自动捕获 print() 语句并作为日志记录 | @flow(log_prints=True)
def demo_flow():
print("Hello from flow") | 若为 False,print 不会出现在日志中。 |
| 缓存配置 | @flow(persist_result=False) | 控制是否持久化 Flow 的返回值 | @flow(persist_result=False)
def temp_calc():
return 42 | 默认 True,可节省存储。 |
| 重试机制 | @flow(retries=3, retry_delay_seconds=10) | 配置 Flow 级重试策略 | @flow(retries=2, retry_delay_seconds=5)
def fragile_flow():
... | 适用于整体流程不稳定场景。 |
3.2 使用 @task 装饰器定义 Task
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
@task 基础用法 | @task
def my_task(x):
return x * 2 | 将函数转换为可追踪的执行单元 | from prefect import task
@task
def add_one(n):
return n + 1 | Task 是最小执行单位。 |
| 命名 Task | @task(name="Process Item") | 自定义 Task 显示名称 | @task(name="Download File")
def download(url):
... | 提高 UI 可读性。 |
| 启用日志捕获 | @task(log_prints=True) | 捕获 print() 输出并记录到日志系统 | @task(log_prints=True)
def debug_task():
print("Debug info") | 推荐开启便于调试。 |
| 结果持久化控制 | @task(persist_result=False) | 不将任务结果写入存储 | @task(persist_result=False)
def temp_op():
return "transient" | 减少 I/O,适合临时数据。 |
| 重试策略 | @task(retries=3, retry_delay_seconds=5) | 设置任务失败后自动重试 | @task(retries=2, retry_delay_seconds=10)
def flaky_api_call():
... | 常用于网络请求类任务。 |
| 缓存机制 | @task(cache_key_fn=..., cache_expiration=...) | 启用结果缓存,避免重复计算 | from prefect.tasks import task_input_hash
@task(cache_key_fn=task_input_hash, cache_expiration=timedelta(hours=1))
def expensive_task(data):
... | task_input_hash 基于输入参数生成缓存键。 |
3.3 Flow 与 Task 的调用与嵌套
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 同步调用 Task | result = my_task() | 在 Flow 中直接调用 Task,等待其完成 | @flow
def parent_flow():
x = add_one(1)
y = add_one(x)
return y | 最常见方式,形成隐式依赖。 |
| 嵌套 Flow | inner_result = child_flow() | 在一个 Flow 中调用另一个 Flow | @flow
def child():
return 10
@flow
def parent():
res = child()
return res * 2 | 子 Flow 也会出现在 UI 中。 |
| 异常传播 | try: ... except: | 捕获 Task 执行中的异常 | @flow
def safe_flow():
try:
risky_task()
except Exception as e:
print(f"Error: {e}") | 未捕获异常会导致 Flow 失败。 |
3.4 参数化 Flow:传递输入参数
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 定义参数化 Flow | @flow
def my_flow(param: str = "default") | 允许 Flow 接收运行时输入 | @flow
def greet(name: str = "World"):
print(f"Hello, {name}!") | 支持类型注解和默认值。 |
| 运行时传参 | my_flow("Alice") 或 CLI 传参 | 调用时传入实际参数 | if __name__ == "__main__":
greet("Bob") | 可通过 Deployment 配置默认参数。 |
| 复杂类型参数 | @flow
def process_data(config: dict) | 接收字典、列表等结构化参数 | @flow
def run_job(settings: dict):
batch_size = settings.get("batch") | 建议使用 Pydantic 模型增强验证。 |
3.5 返回值与数据流传递机制
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| Task 返回值 | return value | 任务执行完成后返回数据 | @task
def get_number():
return 42
@flow
def main():
num = get_number()
print(num) | 返回值可用于后续任务输入。 |
| Flow 返回值 | return final_result | Flow 可返回最终结果供外部使用 | @flow
def calc():
a = task_a()
b = task_b(a)
return a + b | 可用于链式调用或测试。 |
| 数据自动传递 | x = task1()
y = task2(x) | 前一个任务输出作为后一个输入 | @task
def step1(): return "data"
@task
def step2(inp): return inp.upper()
@flow
def pipeline():
d = step1()
result = step2(d) | 构成隐式依赖关系图。 |
第四章:任务执行控制与依赖管理
理解任务间的依赖关系,掌握显式和隐式依赖的构建方式。
4.1 任务依赖的自动推导(隐式依赖)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 基于数据流的依赖 | b = task_b(a) | 当任务 B 使用任务 A 的输出时,自动建立 A → B 依赖 | @flow
def implicit_dep():
x = task1()
y = task2(x) # task2 依赖 task1 | 最自然的依赖方式,无需显式声明。 |
| 多输入依赖 | z = combine(x, y) | 多个任务输出作为同一任务输入 | a = task_a()
b = task_b()
c = combine(a, b) | 所有上游任务完成后才执行 combine。 |
4.2 使用 .submit() 和 .result() 实现异步任务提交(Future 模式)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
.submit() 提交任务 | future = my_task.submit() | 异步提交任务,立即返回 Future 对象 | from prefect import task, flow
@task
def load_data(): return [1,2,3]
@flow
def async_flow():
fut = load_data.submit()
print("Task submitted")
data = fut.result() | 实现非阻塞提交,适合并行。 |
.result() 获取结果 | value = future.result() | 阻塞等待任务完成并获取结果 | 同上 | 若任务未完成会阻塞,可设 timeout。 |
| 并行执行多个任务 | [t.submit() for t in tasks] | 批量提交任务实现并行 | futures = [process_item.submit(i) for i in range(5)]
results = [f.result() for f in futures] | 显著提升吞吐量。 |
4.3 显式依赖:upstream_tasks 与 wait_for
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
wait_for 参数 | task2(wait_for=[task1_future]) | 强制任务等待另一个任务完成,即使无数据传递 | fut = task1.submit()
task2(wait_for=[fut]) | 用于控制执行顺序,无数据依赖时使用。 |
upstream_tasks | @task(upstream_tasks=[task1]) | 装饰器级声明依赖(较少用) | @task(upstream_tasks=[task1])
def task2(): ... | 推荐使用 wait_for 更灵活。 |
4.4 动态任务生成(Mapped Tasks)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
.map() 方法 | task.map(inputs) | 对输入列表中的每个元素并行执行任务 | @task
def square(x): return x ** 2
@flow
def map_flow():
results = square.map([1,2,3,4])
print(results) | 自动并行化,适合批处理。 |
| 多输入映射 | task.map(a_list, b_list) | 多个列表按元素位置映射 | @task
def add(x, y): return x + y
add.map([1,2], [3,4]) → [4, 6] | 列表长度需一致。 |
与 .submit() 结合 | task.submit().map(...) | 高级用法,灵活控制 | 不推荐混合使用,优先使用 .map() | .map() 更简洁安全。 |
4.5 条件分支与跳过任务(if 逻辑与 raise SKIP)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| Python if 逻辑 | if condition: task() | 在 Flow 中使用条件判断控制执行路径 | @flow
def conditional_flow(run_task: bool):
if run_task:
my_task() | 简单分支推荐此方式。 |
| 抛出 SKIP 状态 | from prefect import task, skip
raise skip(reason="...") | 主动跳过任务,标记为 Skipped | @task
def check_file():
if not file_exists():
raise skip("File not found") | 需导入 skip,不视为失败。 |
| 条件性提交 | if condition: task.submit() | 动态决定是否提交任务 | if need_processing:
process_data.submit() | 适用于运行时决策。 |
第五章:状态(State)与结果处理
深入理解任务执行过程中的各种状态及其含义,掌握结果存储策略。
5.1 Prefect 中的 State 类型详解(RUNNING, COMPLETED, FAILED, etc.)
| State 类型 | 说明 | 注意事项 |
|---|
| PENDING | 任务已创建但尚未开始执行 | 通常为初始状态 |
| RUNNING | 任务正在执行中 | 状态会持续到完成或失败 |
| COMPLETED | 任务成功执行并返回有效结果 | 最期望的状态 |
| FAILED | 任务执行过程中抛出未捕获异常 | Flow 中其他任务可能继续执行,取决于依赖 |
| CRASHED | 执行器意外终止(如进程崩溃) | 表示非正常退出 |
| CANCELLED | 用户主动取消任务运行 | 不再继续执行 |
| CANCELLING | 正在处理取消请求的中间状态 | 过渡状态 |
| SCHEDULED | 任务已安排在未来某个时间执行 | 常用于延迟或重试任务 |
| PAUSED | 流程被暂停等待外部事件 | 适用于人工审批等场景 |
| SKIPPED | 任务被跳过(如条件不满足或 raise skip()) | 不视为失败,后续依赖任务可继续 |
5.2 获取任务执行状态对象
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
从 .submit() 获取 Future | future = task.submit()
state = future.wait() | 提交异步任务后等待并获取其状态 | from prefect import task, flow
@task
def my_task(): return 42
@flow
def get_state_flow():
fut = my_task.submit()
state = fut.wait()
print(state.type) # 如 "COMPLETED" | wait() 阻塞直到任务完成 |
使用 .result(timeout=...) | try:
value = future.result(timeout=10)
except TimeoutError:
... | 获取结果的同时获取状态信息 | try:
res = fut.result(timeout=5)
except TimeoutError:
print("Task timed out") | 超时抛出异常,可用于监控 |
| 在 Flow 中捕获 Task 状态 | from prefect import task, allow_failure
result = failed_task.submit().result() | 结合 allow_failure() 使用,获取失败任务的状态 | @task
def may_fail():
raise ValueError("oops")
@flow
def capture_failure():
fut = may_fail.submit()
state = fut.wait()
if state.is_failed():
print("Task failed:", state.message) | 用于容错处理逻辑 |
5.3 自定义状态处理逻辑
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 检查状态类型 | if state.is_completed(): ...
if state.is_failed(): ... | 根据任务执行结果决定后续流程 | state = my_task.submit().wait()
if state.is_completed():
next_step()
elif state.is_failed():
fallback_action() | 推荐使用 is_*() 方法而非直接比较字符串 |
| 自定义失败处理 | try:
result = task.result()
except Exception as e:
handle_error(e) | 捕获异常并执行补偿逻辑 | fut = risky_task.submit()
try:
data = fut.result()
except Exception:
send_alert()
data = get_cached_data() | 适用于高可用场景 |
| 主动设置状态(高级) | 不推荐直接设置 | 自定义任务中可通过 return State(...) 控制,但非常规用法 | 一般通过抛出异常或返回值间接控制状态 | 仅限高级用户扩展使用 |
5.4 结果持久化:启用与配置 Result Storage
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 启用结果持久化 | @task(persist_result=True)
@flow(persist_result=True) | 开启任务/流程结果的自动存储 | @task(persist_result=True)
def expensive_op():
return heavy_calc() | 默认为 True,可关闭以节省 I/O |
| 全局配置结果存储 | prefect config set PREFECT_RESULTS_PERSISTENT=True | 在配置中启用持久化结果 | prefect config set PREFECT_RESULTS_PERSISTENT=True | 影响所有 Flow 和 Task |
| 配置存储路径 | PREFECT_RESULTS_STORAGE_PATH="s3://bucket/results" | 设置结果存储的默认路径 | prefect config set PREFECT_RESULTS_STORAGE_PATH="results/{flow_name}" | 支持模板变量如 {flow_name}, {date} |
| 禁用结果存储 | @task(persist_result=False) | 对特定任务禁用存储 | @task(persist_result=False)
def temp_task():
return "transient" | 适合临时或敏感数据 |
5.5 使用内置 Result 类型(LocalResult, S3Result 等)
| Result 类型 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| LocalResult | 自动使用本地文件系统 | 默认结果存储方式,保存在 .prefect/results | 无需显式配置 | 仅适合本地开发,不适用于分布式环境 |
| S3Result | from prefect_aws import AwsCredentials
from prefect_aws.s3 import S3Bucket
result_storage = S3Bucket.load("my-s3-block") | 存储结果到 AWS S3 | aws_creds = AwsCredentials.load("prod-creds")
s3_block = S3Bucket(
bucket_name="my-data",
credentials=aws_creds
)
@task(result_storage=s3_block)
def upload_data(): ... | 需安装 prefect-aws,配置 Block |
| GCSResult | from prefect_gcp import GcpCredentials
from prefect_gcp.gcs import GCSBucket | 存储到 Google Cloud Storage | 类似 S3 配置方式 | 需安装 prefect-gcp |
| AzureResult | from prefect_azure import AzureBlobStorage | 存储到 Azure Blob Storage | 使用 AzureBlobStorage Block | 需安装 prefect-azure |
| 自定义路径格式 | result_serializer=pickle_serializer | 控制结果序列化方式(如 pickle, json) | from prefect.serializers import PickleSerializer
PickleSerializer() | 注意安全性和跨平台兼容性 |
💡 提示:使用 Block 管理云存储凭证,避免硬编码。
第六章:日志记录与调试技巧
提高可观测性,便于排查问题和监控执行过程。
6.1 在 Task 中使用标准日志(logging)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 配置日志记录器 | import logging
logger = logging.getLogger(__name__) | 创建模块级日志器 | import logging
logger = logging.getLogger("my_module")
@task
def process_item(item):
logger.info(f"Processing {item}") | 推荐使用模块名避免冲突 |
| 输出不同级别日志 | logger.debug(), logger.info(), logger.warning(), logger.error() | 记录不同严重级别的信息 | logger.info("Task started")
logger.warning("Low memory")
logger.error("Failed to connect") | Orion UI 中可按级别过滤 |
| 结构化日志(推荐) | logger.info("Event", extra={"user": "alice", "action": "login"}) | 添加结构化上下文 | logger.info("Data processed", extra={"count": 100, "duration": 2.5}) | 便于后续分析和告警 |
6.2 Prefect 内置日志捕获机制
| 机制 | 说明 | 注意事项 |
|---|
自动捕获 print() | 当 log_prints=True 时,print() 输出会被重定向为日志条目 | 必须在 @flow 或 @task 中启用 |
| 日志级别控制 | 通过 PREFECT_LOGGING_LEVEL=DEBUG 配置全局日志级别 | 支持 DEBUG, INFO, WARNING, ERROR |
| 日志结构化 | Prefect 自动添加任务名、运行 ID、时间戳等元数据 | 无需手动添加上下文 |
| 日志聚合 | 所有日志通过 Orion API 聚合,可在 UI 中查看完整流 | 支持多 Worker 环境 |
6.3 查看运行时日志输出
| 方法 | 用途 | 注意事项 |
|---|
| CLI 运行时输出 | 直接在终端查看 flow() 调用的日志 | 适合本地调试 |
| Orion UI 日志面板 | 访问 http://127.0.0.1:4200 查看 Flow Run 详情页日志 | 支持搜索、过滤、高亮 |
| 日志时间线 | 查看任务执行与日志的时间对应关系 | 用于性能分析 |
| 导出日志 | 通过 API 或 UI 导出日志用于审计 | prefect logs tail 可实时查看 |
6.4 使用 print() 与 logger.info() 的区别
| 对比项 | print() | logger.info() | 注意事项 |
|---|
| 是否被捕获 | 仅当 log_prints=True 时被捕获 | 总是被捕获并结构化 | 推荐统一使用 logging |
| 日志级别 | 无级别,统一视为 INFO | 可设置 DEBUG, WARNING 等 | logging 更灵活 |
| 结构化支持 | 无 | 支持 extra 字段添加上下文 | 便于机器解析 |
| 性能影响 | 轻量 | 可配置异步处理 | 生产环境建议使用 logging |
| 建议用法 | 快速调试 | 生产环境日志记录 | 统一使用 logging 更规范 |
6.5 调试模式运行 Flow(serve, debug 参数)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 启用调试日志 | prefect config set PREFECT_LOGGING_LEVEL=DEBUG | 显示详细内部执行日志 | prefect config set PREFECT_LOGGING_LEVEL=DEBUG
python my_flow.py | 输出大量信息,仅调试时开启 |
使用 serve() 开发模式 | flow.serve(name="dev-deployment") | 持续监听代码变更并自动重启 | @flow
def my_flow(): ...
if __name__ == "__main__":
my_flow.serve(name="test-flow") | Prefect 2.7+ 支持,适合本地开发 |
| 捕获异常并进入调试器 | import pdb; pdb.set_trace() | 在任务中插入断点 | @task
def debug_task():
pdb.set_trace()
return "paused" | 仅限本地,避免提交到生产 |
使用 breakpoint()(Python 3.7+) | breakpoint() | 现代化调试断点 | 同上 | 更简洁,支持 IDE 集成 |
| 禁用并发调试 | @flow(task_runner=SequentialTaskRunner()) | 强制串行执行,便于跟踪 | @flow(task_runner=SequentialTaskRunner())
def debug_flow(): ... | 避免异步干扰调试流程 |
💡 提示:生产环境应关闭 DEBUG 日志,避免性能下降。
第七章:错误处理与重试机制
构建健壮的工作流,应对临时故障。
7.1 异常传播与 Task 失败行为
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 抛出异常导致失败 | raise Exception("error") | 任务中抛出未捕获异常将标记为 FAILED | @task
def failing_task():
raise ValueError("Something went wrong")
@flow
def my_flow():
failing_task() | 异常会中断当前任务执行流 |
| 异常向上游传播 | 默认行为 | 如果任务失败且无处理逻辑,Flow 将进入失败状态 | @flow
def parent_flow():
step1()
failing_task() # 此处失败导致整个 flow 失败 | 后续任务不会执行(除非显式跳过依赖) |
| 依赖中断机制 | 基于数据流依赖 | 若任务 A 失败,依赖其输出的任务 B 不会执行 | a = task_a()
b = task_b(a) # 若 a 失败,b 不运行 | 防止无效或错误数据传播 |
7.2 使用 retry_delay_seconds 和 retries 实现自动重试
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 配置任务级重试 | @task(retries=3, retry_delay_seconds=5) | 失败后自动重试指定次数,每次间隔固定时间 | @task(retries=2, retry_delay_seconds=10)
def unstable_api_call():
response = requests.get("...")
if response.status_code != 200:
raise Exception("Request failed") | 适用于网络抖动、临时服务不可用等场景 |
| Flow 级重试 | @flow(retries=1, retry_delay_seconds=30) | 整个 Flow 执行失败后重试 | @flow(retries=1, retry_delay_seconds=30)
def fragile_pipeline():
task1()
task2() | 慎用,可能导致重复副作用(如写数据库) |
| 指数退避模拟 | 手动实现 | Prefect 不直接支持指数退避,需自定义逻辑 | 结合 time.sleep(random.expovariate(...)) 在任务内部实现 | 推荐用于高频率调用的外部 API |
7.3 自定义重试条件(配合 RetryPolicy)
⚠️ 说明:Prefect 2.x 当前主要通过 retries + 异常捕获实现重试控制,原生 RetryPolicy 更多用于高级扩展。以下为基于最佳实践的”类 RetryPolicy”行为模拟。
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 条件性重试逻辑 | 在任务内捕获特定异常并循环 | 仅对某些异常类型进行重试 | import time
from requests.exceptions import ConnectionError
@task
def retry_on_network_only(data):
for i in range(3):
try:
call_external_service(data)
return
except ConnectionError:
if i == 2: raise
time.sleep(2 ** i)
except Exception as e:
raise e # 其他异常立即失败 | 可精确控制重试条件,但失去 Prefect 内置追踪 |
使用 retry_if 辅助函数(模式) | 自定义判断函数 | 模拟 retry_if 行为 | 定义 should_retry(exception) 函数,在循环中调用 | 属于设计模式,非 Prefect 原生 API |
7.4 捕获特定异常并继续执行
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| try-except 捕获异常 | try: ... except SpecificError: ... | 处理可恢复错误,避免任务失败 | @task
def graceful_task():
try:
result = risky_operation()
except FileNotFoundError:
result = use_default_data()
return result | 提升流程健壮性 |
| 返回默认值 | except: return default | 错误时返回安全默认值 | except KeyError:
return {"status": "unknown"} | 避免下游因空值崩溃 |
| 记录警告而非失败 | except Exception as e:
logger.warning(f"Failed but continuing: {e}") | 继续执行同时记录问题 | 同上 | 适用于非关键路径任务 |
7.5 使用 allow_failure() 忽略某些任务失败
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
allow_failure() 包装器 | from prefect import allow_failure
result = allow_failure(my_task)() | 允许任务失败而不导致 Flow 中断 | from prefect import allow_failure
@task
def may_fail():
raise RuntimeError()
@flow
def tolerant_flow():
optional_result = allow_failure(may_fail)()
print("Continuing despite failure")
next_step() | 失败任务返回 None |
与 .submit() 结合使用 | future = allow_failure(task).submit() | 异步提交并容忍失败 | fut = allow_failure(risky_task).submit()
state = fut.wait()
if state.is_failed():
handle_gracefully() | 仍可获取状态进行处理 |
| 适用场景 | —— | 日志清理、通知发送、非关键检查等 | # 发送 Slack 通知,失败也不影响主流程
notify_status = allow_failure(send_slack).submit() | 明确标识哪些任务是”尽力而为” |
第八章:任务运行模式与并发控制
控制任务是串行、并行还是异步执行,优化性能。
8.1 默认同步执行模型
| 特性 | 说明 | 注意事项 |
|---|
| 执行顺序 | 任务按代码顺序依次执行,前一个完成后再开始下一个 | 最简单直观的行为 |
| 并发性 | 单线程内串行执行,不利用多核或多线程 | I/O 密集型任务效率低 |
| 示例行为 | task_a() → 等待完成 → task_b() → 等待完成 → task_c() | 适合轻量或有强顺序依赖的流程 |
| 适用场景 | 本地调试、小型脚本、强依赖链 | 初学者默认模式 |
8.2 启用异步 Task 支持(async/await)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 定义异步 Task | @task
async def async_task():
await asyncio.sleep(1) | 创建可挂起的异步任务 | import asyncio
from prefect import task, flow
@task
async def fetch_data(url):
await asyncio.sleep(0.5)
return f"Data from {url}" | 需使用 async def |
| 调用异步 Task | await async_task() | 在异步 Flow 中调用 | @flow
async def async_flow():
result = await fetch_data("test.com")
return result | Flow 也必须是 async |
| 并发异步任务 | results = await asyncio.gather(t1(), t2()) | 同时运行多个异步任务 | urls = ["a.com", "b.com"]
tasks = [fetch_data(url) for url in urls]
results = await asyncio.gather(*tasks) | 极大提升 I/O 密集型性能 |
8.3 使用 concurrency_limit 限制并行度
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 设置并发限制 | @task(concurrency_limit=3) | 控制同名任务的最大并发数 | @task(concurrency_limit=2)
def limited_task(n):
time.sleep(2)
return n * 2
@flow
def limit_flow():
futures = [limited_task.submit(i) for i in range(5)]
[f.result() for f in futures] | 同一时刻最多 2 个 limited_task 运行 |
| 全局配置 | prefect config set PREFECT_TASKS_DEFAULT_CONCURRENCY_LIMIT=5 | 设置所有任务的默认并发限制 | prefect config set PREFECT_TASKS_DEFAULT_CONCURRENCY_LIMIT=1 | 影响未显式设置的任务 |
| 动态调整 | 通过 UI 或 API 修改 | 运行时动态控制资源使用 | 暂不支持直接代码修改,需通过 Block 或部署配置 | 适合弹性调度场景 |
8.4 使用不同的 Task Runner
| Task Runner 类型 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| SequentialTaskRunner | task_runner=SequentialTaskRunner() | 强制串行执行,一个接一个 | @flow(task_runner=SequentialTaskRunner())
def serial_flow(): ... | 调试时确保执行顺序 |
| ConcurrentTaskRunner | task_runner=ConcurrentTaskRunner() | 默认,使用线程池并发执行可等待任务 | @flow(task_runner=ConcurrentTaskRunner())
def concurrent_flow(): ... | 适合 I/O 密集型任务 |
| ThreadPoolTaskRunner | from prefect.task_runners import ThreadPoolTaskRunner
task_runner=ThreadPoolTaskRunner(max_workers=4) | 自定义线程池大小 | runner = ThreadPoolTaskRunner(max_workers=8)
@flow(task_runner=runner)
def threaded_flow(): ... | 控制并发线程数 |
| ProcessPoolTaskRunner | from prefect.task_runners import ProcessPoolTaskRunner
task_runner=ProcessPoolTaskRunner() | 使用多进程并行,适合 CPU 密集型任务 | @flow(task_runner=ProcessPoolTaskRunner())
def cpu_intensive_flow(): ... | 避免 GIL 限制,但进程开销大 |
| 设置位置 | 传递给 @flow 装饰器 | Flow 级别控制执行策略 | 所有示例均在 @flow 中设置 task_runner 参数 | 不能在 Task 上单独设置 |
8.5 异步 Flow 中运行同步 Task
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 直接调用同步 Task | result = sync_task() | 在 async Flow 中调用普通 @task | @task
def sync_op(x): return x + 1
@flow
async def mixed_flow():
a = await async_task()
b = sync_op(a) # 直接调用 | 同步任务会阻塞事件循环 |
| 提交同步 Task 异步运行 | future = sync_task.submit()
result = await future.result() | 避免阻塞,让 Task Runner 处理 | fut = sync_task.submit(a)
result = await fut.result() | 推荐方式,提升并发性 |
| 混合模式最佳实践 | 结合使用 | 构建高效混合流水线 | @flow
async def best_practice_flow():
data = await fetch_data_async()
processed = await run_cpu_task_in_process_pool(data)
await log_to_db.submit(processed)
return processed | 根据任务类型选择执行方式 |
💡 提示:在异步 Flow 中运行大量同步任务时,建议使用 ProcessPoolTaskRunner 避免阻塞事件循环。
第九章:部署与调度(Deployments)
将本地 Flow 部署为可长期运行的服务,支持定时触发。
9.1 什么是 Deployment?
| 概念 | 说明 | 注意事项 |
|---|
| Deployment | 一个可调度的 Flow 实例,包含运行配置(参数、调度、基础设施等) | 是将本地 Flow 发布为生产服务的关键机制 |
| 核心组成 | Flow 函数、参数默认值、调度计划(可选)、基础设施配置(如 Docker、Kubernetes)、存储位置(如 S3、GitHub) | 所有信息可定义在 YAML 文件中 |
| 与 Flow 的区别 | Flow 是代码逻辑,Deployment 是运行时配置 | 同一个 Flow 可有多个 Deployment(如 dev/prod) |
| 生命周期 | 长期存在,可被调度器自动触发或手动运行 | 不依赖本地机器运行 |
| 使用场景 | 定时任务、事件驱动流水线、API 触发任务等 | 替代 crontab 的现代化调度方案 |
9.2 创建 Deployment 配置文件(YAML)
| 字段 | 语法 | 用途 | 示例值 | 注意事项 |
|---|
name | name: "daily-ingest" | 部署名称,用于区分不同环境 | "nightly-sync" | 建议包含环境或用途 |
flow_name | flow_name: "data_pipeline" | 对应的 Flow 函数名 | "etl_flow" | 必须与代码中 @flow(name=...) 或函数名一致 |
parameters | parameters:
config: "prod"
batch_size: 1000 | 设置 Flow 运行时默认参数 | 可为空 | 支持复杂结构如 dict/list |
schedule | schedule:
cron: "0 2 * * *"
timezone: "UTC" | 定义自动调度规则 | "0 0 * * MON"(每周一) | 可选,无则仅支持手动触发 |
work_pool | work_pool: "k8s-prod" | 指定执行该部署的工作池 | "docker-dev" | 必须已注册对应 Work Pool |
entrypoint | entrypoint: "flows/etl.py:data_pipeline" | 指向 Flow 模块路径 | "src/main.py:my_flow" | 格式:<file_path>:<flow_name> |
infra_overrides | infra_overrides:
env.PREFECT_LOGGING_LEVEL: DEBUG | 覆盖基础设施环境变量 | 用于临时调试 | 生产环境慎用 |
📄 示例 deployment.yaml:
name: daily-etl-prod
flow_name: etl_flow
entrypoint: flows/etl.py:etl_flow
parameters:
environment: "production"
batch_size: 5000
schedule:
cron: "0 3 * * *"
timezone: "Asia/Shanghai"
work_pool: aws-ecs-prod
9.3 使用 prefect deploy 命令打包和上传
| 命令 | 语法 | 用途 | 示例 | 注意事项 |
|---|
| 初始化部署配置 | prefect deployment build <entrypoint> -n <name> | 生成 deployment.yaml 模板 | prefect deployment build ./my_flow.py:main_flow -n dev-run | 自动生成基础配置文件 |
| 添加调度 | -s "cron_expression" | 在构建时添加调度 | prefect deployment build ... -s "0 2 * * *" | 等价于在 YAML 中写 schedule |
| 指定工作池 | -q <queue_name> | 指定目标工作池 | prefect deployment build ... -q k8s-prod | 必须已存在该 Work Pool |
| 构建并应用 | --apply | 构建后立即上传到 Orion | prefect deployment build ... --apply | 一键完成部署注册 |
| 应用现有配置 | prefect deploy -n <deployment_name> | 根据本地 YAML 文件部署 | prefect deploy -n daily-etl-prod | 需先有 deployment.yaml |
| 查看所有部署 | prefect deployment ls | 列出已注册的 Deployments | prefect deployment ls | 检查状态是否为 READY |
9.4 设置 Cron 调度计划(CronSchedule)
| 方法 | 语法 | 用途 | 示例 | 注意事项 |
|---|
| CLI 构建时设置 | -s "0 0 * * *" --timezone "UTC" | 在 build 命令中添加调度 | prefect deployment build ... -s "0 15 * * MON-FRI" | 支持标准 cron 表达式 |
| YAML 配置 | schedule:
cron: "*/30 * * * *"
timezone: "America/New_York" | 在 deployment.yaml 中声明 | 使用 crontab.guru 验证表达式 | 时区建议明确设置 |
| Python 中定义 | from prefect.deployments import Deployment
from prefect.schedules import CronSchedule
Deployment.build_from_flow(
flow=my_flow,
name="nightly",
schedule=CronSchedule(cron="0 2 * * *")
) | 代码方式创建调度 | 适合 CI/CD 自动化 | 需调用 .apply() 注册 |
| 常见表达式 | —— | —— | "0 * * * *"(每小时)
"0 0 1 * *"(每月1号)
"@daily"(每天0点) | Prefect 支持 @hourly, @daily 等简写 |
9.5 手动触发与暂停 Deployment
| 操作 | 命令/方式 | 用途 | 示例 | 注意事项 |
|---|
| 手动触发运行 | prefect deployment run <deployment_name> | 立即执行一次部署 | prefect deployment run data-sync-prod | 不受调度限制,可传额外参数 |
| 传参运行 | --param key=value | 覆盖默认参数 | prefect deployment run etl-flow --param date=2025-01-01 | 支持多个 --param |
| 暂停部署 | prefect deployment pause <name> | 停止自动调度,但仍可手动触发 | prefect deployment pause nightly-backup | 避免误执行 |
| 恢复部署 | prefect deployment resume <name> | 重新启用自动调度 | prefect deployment resume nightly-backup | 暂停期间错过的调度不会补发 |
| UI 操作 | Orion UI → Deployments → Actions | 图形化触发/暂停 | 点击 “Run” 或 “Pause” 按钮 | 更直观,适合非技术人员 |
9.6 多环境部署管理(dev/staging/prod)
| 策略 | 实现方式 | 示例 | 注意事项 |
|---|
| 不同 Deployment 名称 | my-flow-dev, my-flow-staging, my-flow-prod | 通过命名区分环境 | 配合 CI/CD 脚本自动部署 |
| 不同 Work Pool | dev-docker, staging-k8s, prod-ecs | 隔离执行资源 | 避免开发任务占用生产资源 |
| 参数区分环境 | parameters:
env: "dev" vs env: "prod" | Flow 内部根据参数切换逻辑 | if params["env"] == "prod": send_alert() |
| 多个 YAML 文件 | deployment-dev.yaml, deployment-prod.yaml | 分别配置不同环境 | prefect deploy -f deployment-prod.yaml |
| CI/CD 集成 | GitHub Actions / GitLab CI | 自动构建并部署到指定环境 | PR 合并 → 部署到 staging 手动批准 → 部署到 prod |
第十章:Prefect Cloud / Orion UI 使用指南
利用可视化界面监控、管理和分析工作流执行。
10.1 登录与项目初始化(Prefect Cloud)
| 功能 | 操作方式 | 说明 | 注意事项 |
|---|
| 访问 Orion UI | 浏览器打开 http://127.0.0.1:4200(本地)或 https://app.prefect.cloud(Cloud) | 可视化控制台入口 | 本地 Orion 需运行 prefect orion start |
| 登录 Prefect Cloud | 访问 https://app.prefect.cloud | 使用 GitHub 或邮箱注册/登录 | 提供免费 tier 和团队协作功能 |
| 创建 Workspace | Cloud 中点击 “Create Workspace” | 隔离不同项目或团队的数据 | 如 analytics-prod, ml-training |
| 初始化本地配置 | prefect cloud login -k <api_key> -w <workspace> | 将本地 CLI 连接到 Cloud | API Key 在用户设置中生成 |
| 查看连接状态 | prefect whoami | 确认当前登录用户和工作区 | 检查是否连接正确 |
10.2 查看 Flows 与 Deployments 列表
| 视图 | 位置 | 内容 | 注意事项 |
|---|
| Flows 列表 | 左侧菜单 → Flows | 显示所有已注册的 Flow | 按名称、创建时间排序 |
| Deployment 列表 | 左侧菜单 → Deployments | 显示所有部署及其调度状态 | 可见 SCHEDULED, PAUSED, READY 状态 |
| 搜索功能 | 页面顶部搜索框 | 按名称过滤 Flows/Deployments | 支持模糊匹配 |
| 状态标识 | 图标颜色 | 绿色:正常,黄色:警告,红色:失败 | 快速识别问题 |
| 版本管理 | (未来版本) | 当前版本不强制版本化,依赖 Git | 建议结合 Git 提交记录追踪 |
10.3 监控 Flow Runs 与 Task Runs 状态
| 功能 | 位置 | 说明 | 注意事项 |
|---|
| Flow Runs 列表 | 点击某个 Flow → “Runs” 标签页 | 查看该 Flow 的所有执行记录 | 按时间倒序排列 |
| Task Runs 详情 | 点击某次 Flow Run → 下钻 | 查看每个 Task 的执行状态 | 可见 COMPLETED, FAILED, SKIPPED 等 |
| 状态时间线 | 时间轴视图 | 可视化任务启动/结束时间 | 分析性能瓶颈 |
| 过滤器 | 状态、时间范围、标签等 | 筛选特定运行记录 | 如只看 FAILED 的运行 |
| 重试按钮 | 在失败的 Flow Run 上 | 重新执行该次运行 | 不会创建新 Run,更新原记录 |
10.4 查阅日志与时间线视图
| 功能 | 位置 | 说明 | 注意事项 |
|---|
| 实时日志流 | Flow Run 详情页 → “Logs” 标签 | 查看任务打印的日志和 print() 输出 | 支持滚动加载 |
| 日志搜索 | 日志面板内搜索框 | 按关键字查找日志条目 | 如搜索 “error” |
| 时间线视图 | ”Timeline” 标签页 | 图形化展示任务执行时间、依赖关系 | 识别并行/串行模式 |
| 任务详情弹窗 | 点击某个 Task Run | 查看开始时间、持续时间、返回值、状态 | 用于性能分析 |
| 结构化日志 | 日志条目包含 task_name, flow_run_id 等字段 | 便于机器解析和审计 | 推荐使用 logging 而非 print |
10.5 设置通知(Notifications)与告警规则
| 功能 | 配置方式 | 支持渠道 | 注意事项 |
|---|
| 创建通知块(Block) | UI → Blocks → + New → Notification | Slack Webhook、Email(SMTP)、Microsoft Teams、PagerDuty、Custom Webhook | 需提供 URL 或凭证 |
| 绑定告警规则 | 在 Flow 或 Deployment 设置中 | 选择触发条件和通知目标 | 如 “Flow Run Failed” → Slack Channel |
| 常见触发条件 | Flow Run Failed、Flow Run Crashed、Flow Run Late(未按时启动)、State Changed | 可自定义 | 推荐为关键流程设置失败通知 |
| 测试通知 | Block 页面提供 “Test” 按钮 | 发送测试消息验证配置 | 部署前务必测试 |
| 多通知组合 | 可为同一事件配置多个通知 | 如失败时同时发 Slack 和邮件 | 确保关键人员能收到 |
💡 提示:使用 prefect alerts create 命令也可通过 CLI 创建告警规则。
第十一章:高级特性与扩展功能
探索更复杂的应用场景与集成能力。
11.1 使用 Secrets 管理敏感信息(Secret Blocks)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 创建 Secret Block | prefect block create secret my-password | 安全存储密码、API Key 等 | 在 CLI 中交互式输入值 | 值加密存储于数据库 |
| 加载 Secret | from prefect.blocks.core import Secret
secret_block = Secret.load("my-api-key")
api_key = secret_block.get() | 在 Task 中安全获取敏感信息 | @task
def call_api():
key = Secret.load("prod-api-key").get()
requests.get(..., headers={"X-API-Key": key}) | 避免日志泄露 |
| 环境变量注入 | PREFECT_ORION_API_KEY={{secret:my-api-key}} | 在部署或基础设施中引用 Secret | 用于连接外部服务 | 支持模板语法 {{secret:xxx}} |
| UI 管理 | Orion UI → Blocks → + New → Secret | 图形化创建和管理 | 适合非代码用户 | 名称区分环境(如 db-password-prod) |
🔐 安全建议:绝不将敏感信息硬编码或提交至 Git。
11.2 集成外部系统:Blocks(Storage, Infrastructure, Notification)
| Block 类型 | 示例 | 用途 | 配置方式 | 注意事项 |
|---|
| Storage Blocks | GitHub, S3Bucket, AzureBlobStorage, GCSBucket | 存储 Flow 代码或结果 | GitHub(repository="myorg/flows", branch="main") | 用于远程加载 Flow |
| Infrastructure Blocks | DockerContainer, KubernetesJob, ECSTask, Process | 定义任务运行环境 | DockerContainer(image="my-image:latest") | 决定任务在哪执行 |
| Notification Blocks | SlackWebhook, Email SMTP, MicrosoftTeamsWebhook | 发送告警或通知 | SlackWebhook(url="https://hooks.slack.com/...") | 需测试连通性 |
| Credential Blocks | AwsCredentials, GcpCredentials, AzureCredentials | 管理云平台认证 | AwsCredentials(access_key_id=..., secret_access_key=...) | 推荐使用 IAM 角色 |
| 使用方式 | .load("block-name") | 在代码或部署中引用 | infra = DockerContainer.load("prod-container") | Block 必须已注册 |
11.3 自定义 Block 开发
| 步骤 | 语法 | 说明 | 示例 | 注意事项 |
|---|
| 定义 Block 类 | from prefect.blocks.core import Block
class MyCustomBlock(Block):
field1: str
field2: int = 42 | 继承 Block 并定义字段 | 可添加验证逻辑(@validator) | 字段应有类型注解 |
| 注册 Block | MyCustomBlock.save("my-instance", overwrite=True) | 将配置保存为可复用 Block | 在脚本中运行一次注册 | 支持 CLI save |
| 加载使用 | block = MyCustomBlock.load("my-instance") | 在 Flow/Task 中使用 | conn = block.connect() | 可封装连接逻辑 |
| 高级功能 | @property, @secret 字段 | 添加计算属性或标记敏感字段 | api_key: str = Field(..., json_schema_extra={"secret": True}) | @secret 自动加密 |
| 发布为插件 | 打包为 Python 包 | 供团队共享使用 | pip install my-prefect-blocks | 遵循 Prefect 插件规范 |
💡 提示:自定义 Block 可极大提升团队标准化程度。
11.4 使用 Events & Observability API
| 功能 | 语法 | 用途 | 示例 | 注意事项 |
|---|
| 发送自定义事件 | from prefect.events.emit import emit_event
emit_event("my.custom.event", resource={"prefect.resource.id": "my-flow"}) | 记录业务事件,用于监控或告警 | emit_event("data.validation.failed", payload={"file": "data.csv"}) | 事件可被自动化规则捕获 |
| 事件结构 | event, resource, related, payload | 标准化事件格式 | resource={"prefect.resource.id": "etl-flow"} | 符合 CloudEvents 规范 |
| 与 Automations 集成 | 在 Prefect Cloud 中创建 Automation | 响应事件触发动作 | 如 “当 db.backup.failed 事件发生 → 发送 Slack” | 替代传统轮询告警 |
| 查询事件流 | prefect events tail | 实时查看事件流 | 调试或审计用途 | 需权限 |
| 用途场景 | 审计日志、合规报告、外部系统集成 | 超越任务状态的上下文信息 | 与 SIEM 系统对接 | 属于高级可观测性功能 |
11.5 与 CI/CD 流程集成
| 阶段 | 实现方式 | 示例 | 注意事项 |
|---|
| 代码提交 | Git Push 触发 CI(GitHub Actions / GitLab CI) | on: [push] | 仅允许特定分支触发部署 |
| 流程测试 | 运行 pytest 或本地 flow() 验证 | python test_flows.py | 确保逻辑正确 |
| 构建部署 | prefect deployment build ... --apply | 自动生成并注册 Deployment | 使用 --apply 一键发布 |
| 多环境部署 | 条件判断分支或环境变量 | if [ "$ENV" = "prod" ]; then prefect deploy -n prod-flow; fi | 分阶段发布(dev → staging → prod) |
| 版本标记 | 结合 Git Tag 或镜像版本 | image: my-flow:$GIT_SHA | 便于回滚 |
| 安全控制 | 使用 Service Account Token | 避免使用个人 API Key | 最小权限原则 |
🔄 推荐流程:Git Push → CI 测试 → 构建部署到 dev → 手动批准 → 部署到 prod
11.6 使用 Server 模式自托管 Orion
| 配置项 | 方法 | 说明 | 注意事项 |
|---|
| 启动本地 Server | prefect orion start | 默认在 http://127.0.0.1:4200 启动 | 数据存储于 SQLite(默认) |
| 使用 PostgreSQL | prefect config set PREFECT_ORION_DATABASE_CONNECTION_URL=postgresql+asyncpg://... | 生产环境推荐使用 | 支持高并发和持久化 |
| 配置服务器地址 | prefect config set PREFECT_API_URL=http://my-orion:4200/api | 指向自托管实例 | 所有 Agent 和 CLI 使用该地址 |
| 部署 Agent | prefect agent start -q 'default' | 监听部署并执行任务 | Agent 需能访问 API |
| TLS/HTTPS | 反向代理(Nginx, Traefik) | 对外暴露安全接口 | 生产环境必须启用 |
| 备份策略 | 定期备份数据库 | 防止元数据丢失 | PostgreSQL 支持标准备份工具 |
🏢 适用场景:数据合规要求高、内网部署、定制化需求强的企业。
第十二章:最佳实践与性能优化
总结生产级使用经验,避免常见陷阱。
12.1 Flow 设计原则:单一职责、高内聚低耦合
| 原则 | 说明 | 反模式 | 建议 |
|---|
| 单一职责 | 一个 Flow 只做一件事(如 ETL、通知、清理) | 一个 Flow 包含数据清洗、训练、部署模型 | 拆分为 clean_data_flow, train_model_flow |
| 高内聚 | 相关任务组织在同一 Flow | 将用户注册和日志归档混在一个 Flow | 按业务域划分 |
| 低耦合 | Flow 间通过参数或结果传递数据,避免直接依赖内部实现 | Flow A 直接访问 Flow B 的临时文件 | 使用 result_storage 或 API 通信 |
| 可复用性 | 设计通用 Flow,通过参数定制行为 | 每个客户写一个硬编码 Flow | run_pipeline(customer_id="c1") |
12.2 避免在 Flow 层做 heavy 计算
| 问题 | 风险 | 解决方案 | 示例 |
|---|
| Flow 中直接处理大数据 | 阻塞调度器、内存溢出 | 将计算放入 @task | ❌ data = [x*2 for x in large_list] ✅ @task
def process_data(data): ... |
| 复杂逻辑判断 | Flow 变得臃肿难维护 | 抽象为独立 Task | @task
def should_run_full_pipeline(): ... |
| 长时间阻塞 | 影响其他 Flow 调度 | 使用异步或分离计算任务 | 使用 ProcessPoolTaskRunner |
| 建议 | Flow 应像”指挥家”,协调任务而非亲自执行 | —— | Flow 负责依赖、调度、错误处理 |
12.3 合理划分 Task 粒度
| 粒度 | 优点 | 缺点 | 推荐场景 |
|---|
| 过细(太多小 Task) | 高并发、独立重试 | 调度开销大、日志爆炸 | 不推荐 |
| 过粗(一个大 Task) | 减少开销 | 失败需整体重试、无法并行 | CPU 密集型且原子性强 |
| 适中(合理拆分) | 平衡开销与灵活性 | 需设计判断 | —— |
| 推荐做法 | 按功能边界拆分 | —— | fetch_data() → clean_data() → upload_result() |
✅ 黄金法则:一个 Task 应能独立描述其功能。
12.4 使用缓存加速重复任务(cache_key_fn, cache_expiration)
| 参数 | 语法 | 用途 | 示例 | 注意事项 |
|---|
cache_key_fn | @task(cache_key_fn=task_input_hash) | 基于输入生成缓存键 | from prefect.tasks import task_input_hash
@task(cache_key_fn=task_input_hash, cache_expiration=timedelta(hours=1))
def expensive_op(data): ... | 相同输入直接返回缓存结果 |
cache_expiration | cache_expiration=timedelta(days=1) | 设置缓存有效期 | 避免使用过期数据 | —— |
| 自定义缓存键 | cache_key_fn=lambda: "static-key" | 强制使用固定键 | 适用于无参任务 | —— |
| 适用场景 | 幂等、耗时、结果稳定的任务 | API 调用、大文件解析、模型加载 | 非幂等操作(如写数据库)禁用缓存 | —— |
12.5 监控资源消耗与超时设置
| 配置 | 语法 | 用途 | 建议值 | 注意事项 |
|---|
timeout_seconds | @task(timeout_seconds=300) | 防止任务无限挂起 | 根据任务类型设置(如 API 调用 30s) | 超时抛出 TaskTimeoutError |
| 资源监控 | Orion UI → Flow Run → Metrics | 查看内存、CPU 使用 | 结合外部 APM(如 Prometheus) | 识别性能瓶颈 |
| 重试与超时配合 | retries=2, retry_delay_seconds=10 | 应对临时超时 | 避免雪崩(指数退避更佳) | 超时任务可重试 |
| 基础设施资源限制 | DockerContainer(env={"MAX_WORKERS": "4"}) | 控制容器资源 | 与 Kubernetes resources 配合 | 防止单任务耗尽资源 |
12.6 版本管理与 Deployment 升级策略
| 策略 | 实现方式 | 说明 | 注意事项 |
|---|
| Git 作为唯一源 | 所有 Flow 代码存于 Git | 支持审计、回滚、CI/CD | 推荐使用 |
| 镜像版本化 | DockerContainer(image="my-flow:v1.2.3") | 结合 Docker Tag 管理版本 | latest 不推荐用于生产 |
| Deployment 更新 | prefect deploy 覆盖已有部署 | 更新代码或配置后重新应用 | 不影响正在运行的实例 |
| 蓝绿部署 | 维护 flow-v1 和 flow-v2 两个 Deployment | 流量切换前验证新版本 | 需外部调度器支持 |
| 回滚机制 | 重新部署旧版 YAML 或镜像 | 快速恢复服务 | 保留历史镜像和配置 |
💡 建议:每次变更生成新版本,避免就地修改 — 版本 = Git Commit + 镜像 Tag