R G
10/23/2025, 8:53 AMMarvin
10/23/2025, 8:53 AMMarvin
10/23/2025, 9:06 AMmy_flow()) is blocking, so a single thread will run one flow at a time. However, a single Python process can run multiple flows concurrently if you orchestrate them with threads, asyncio, or subprocesses.
Common patterns:
- Threads
from threading import Thread
from prefect import flow
@flow
def my_flow(n: int):
...
t1 = Thread(target=my_flow, args=(1,))
t2 = Thread(target=my_flow, args=(2,))
t1.start(); t2.start()
t1.join(); t2.join()
- Asyncio (with async flows)
import asyncio
from prefect import flow
@flow
async def my_flow(n: int):
...
asyncio.run(asyncio.gather(my_flow(1), my_flow(2)))
- Separate processes for full isolation (advanced): you can spawn subprocesses to run flows if needed for isolation or GIL-bound work.
Notes and caveats:
- Prefect contexts are isolated via contextvars, so running flows in different threads or asyncio tasks is safe; each flow gets its own run context. See Runtime context.
- Within a single flow, task concurrency is controlled by the task runner (e.g., ThreadPoolTaskRunner, ProcessPoolTaskRunner). That’s separate from running multiple flows in parallel. See Task runners.
- If your concurrent flows share global state, add your own locking. For subprocesses, ensure data is picklable.
- How-to guide: Run work concurrently. Global concurrency controls: Global concurrency limits.
If you share a bit about your workload (CPU vs I/O bound, sync vs async, need for isolation), I can suggest the best pattern.R G
10/23/2025, 9:12 AMasyncio.run, is it necessary to use <http://asyncio.to|asyncio.to>_thread when invoking the flows?Marvin
10/23/2025, 9:12 AMMarvin
10/23/2025, 9:16 AM@flow on an async def), just await it. No need for asyncio.to_thread.
- If your flow is sync (def), offload it so you don’t block the event loop. You can use asyncio.to_thread(...) or Prefect’s helper run_sync_in_worker_thread(...).
Examples:
- Async flow
from prefect import flow
@flow
async def process(data):
...
async def handle_request(req):
result = await process(req.data) # no to_thread needed
return result
- Sync flow
import asyncio
from prefect import flow
@flow
def process(data):
...
async def handle_request(req):
result = await asyncio.to_thread(process, req.data)
return result
- Using Prefect’s helper (preferred for sync flows)
from prefect import flow
from prefect.utilities.asyncutils import run_sync_in_worker_thread
@flow
def process(data):
...
async def handle_request(req):
result = await run_sync_in_worker_thread(process, req.data)
return result
Notes:
- Prefect’s contexts are thread- and asyncio-safe via contextvars, so running flows in threads won’t corrupt state. See Runtime context.
- Async utils reference: prefect.utilities.asyncutils.
If you can share whether your flow is sync or async and whether the work is CPU- or I/O-bound, I can recommend the best pattern.R G
10/24/2025, 2:21 AMMarvin
10/24/2025, 2:21 AMMarvin
10/24/2025, 2:26 AMfrom prefect import flow, task
@task(cache_result_in_memory=False, persist_result=False)
def step(...):
...
@flow(cache_result_in_memory=False, persist_result=False, log_prints=False)
def my_flow(...):
...
- Only enable persist_result=True if you need to read results later; otherwise keep it off to reduce serialization and references.
- Docs: Results and result storage, Caching concepts
2) Avoid keeping large objects alive
- Don’t return large DataFrames/arrays from tasks if you can help it; write to disk/object storage and return a small reference (path/URL/ID).
- If you submit many tasks, resolve and drop references promptly:
f = some_task.submit(...)
result = f.result()
del f, result
- If you truly must return a large final result to the caller, ensure intermediate results are not retained.
3) Use process isolation to reclaim memory reliably
- If the growth is driven by underlying libraries, isolate work in processes so memory is released when the process exits.
- Option A: Run tasks in a process pool within the flow:
from prefect import flow
from prefect.task_runners import ProcessPoolTaskRunner
@flow(task_runner=ProcessPoolTaskRunner(), cache_result_in_memory=False)
def my_flow(...):
...
- Option B: Run the entire flow in a subprocess from your service:
from prefect.flow_engine import run_flow_in_subprocess
proc = run_flow_in_subprocess(my_flow, parameters={"x": 1})
proc.join()
# exchange results via files/object storage if needed
- API ref: Task runners, Flow engine API
4) Reduce logging pressure
- Disable print capture if not needed: @flow(log_prints=False)
- Lower log verbosity to avoid large in-memory buffers. See Logging handlers.
5) Coarsen tiny tasks or stream/batch work
- Thousands of tiny tasks can create overhead. Batch small units inside a single task where possible.
- If you’re building up large collections in memory, restructure to process in chunks and free memory between chunks.
A couple of clarifying questions to tailor this further:
- Roughly how many tasks run per flow and what kinds of objects are returned?
- Do you need to keep intermediate task outputs available, or can they be discarded/persisted externally?
- Is the workload CPU-bound (e.g., heavy pandas/NumPy) or I/O-bound?
If you share those details, I can suggest the most effective combination (e.g., disabling in-memory caching + process pool vs. full subprocess isolation) for your case.R G
10/24/2025, 2:41 AMMarvin
10/24/2025, 2:41 AMMarvin
10/24/2025, 2:43 AMfrom prefect import flow, task
from prefect.task_runners import ThreadPoolTaskRunner
@task(cache_result_in_memory=False, persist_result=False)
def io_step(...):
...
@flow(
cache_result_in_memory=False,
persist_result=False,
log_prints=False, # avoid log buffering
task_runner=ThreadPoolTaskRunner(max_workers=8) # good for I/O-bound
)
def my_flow(...):
# chain steps without retaining large intermediates
a = io_step(...)
b = io_step(a)
...
return final_small_result
- Docs: Caching and Results
2) Don’t keep references to intermediates
- Avoid accumulating futures/results in global lists/dicts. Consume results and drop references promptly.
f = io_step.submit(...)
res = f.result()
# use res then
del f, res
- If any task produces very large data, write it to disk/object storage and return a small reference (path/URL/ID) instead of the full object.
3) Use threads (not processes) for I/O-bound work
- Threads are ideal for I/O-bound tasks and avoid the pickling overhead of processes. Use ThreadPoolTaskRunner as shown. Limit max_workers to a sensible number to avoid too many open sockets/files.
4) Reduce logging memory
- Set log_prints=False unless you truly need it, and keep log volume reasonable to avoid buffering overhead.
- Logging reference: Logging handlers
5) As a belt-and-suspenders option: process isolation per run
- If memory still rises run-over-run (e.g., due to third-party libs), run the entire flow in a subprocess so memory is guaranteed to be released when the process exits.
from prefect.flow_engine import run_flow_in_subprocess
# Parent service code:
proc = run_flow_in_subprocess(my_flow, parameters={"...": "..."})
proc.join()
# For large final results, write to a file/object store inside the flow and read back here.
- API: Flow engine API
Extra note for your asyncio service:
- If your flow is sync, call it without blocking the event loop:
from prefect.utilities.asyncutils import run_sync_in_worker_thread
async def handle_request(...):
result = await run_sync_in_worker_thread(my_flow, ...)
return result
- If your flow is async (async def), just await my_flow(...).
If you implement 1–4, memory typically stabilizes. If it still grows across requests, 5 (subprocess per run) provides the strongest isolation.R G
10/24/2025, 3:09 AMrun_flow_in_subprocess specifically and how to obtain the flow's running results.Marvin
10/24/2025, 3:09 AMMarvin
10/24/2025, 3:10 AMrun_flow_in_subprocess gives you memory isolation, but since it runs the flow in a separate process, it doesn’t return the flow’s Python result directly. The usual pattern is:
- Pass the flow a path to write its final result (e.g., JSON, pickle, parquet)
- Start the subprocess
- Wait for it to finish
- Read the result from that path
Below are complete examples you can drop in.
Flow code (write the final result to a provided path)
import json
from prefect import flow, task
@task(cache_result_in_memory=False, persist_result=False)
def io_step(val: int) -> dict:
# do I/O-bound work and produce a small result
return {"value": val}
@flow(cache_result_in_memory=False, persist_result=False, log_prints=False)
def my_flow(x: int, output_path: str) -> None:
# chain your ~20 tasks here; only final result matters
final = io_step(x)
with open(output_path, "w") as f:
json.dump(final, f)
Synchronous service code (run in a subprocess and read the result)
import os
import json
import tempfile
from prefect.flow_engine import run_flow_in_subprocess
def run_my_flow_and_get_result(x: int) -> dict:
with tempfile.TemporaryDirectory() as tmpdir:
output_path = os.path.join(tmpdir, "result.json")
proc = run_flow_in_subprocess(my_flow, parameters={"x": x, "output_path": output_path})
proc.join()
if proc.exitcode != 0:
raise RuntimeError(f"Flow subprocess failed with exit code {proc.exitcode}")
with open(output_path, "r") as f:
return json.load(f)
If your service is async (don’t block the event loop)
import asyncio
from prefect.utilities.asyncutils import run_sync_in_worker_thread
async def run_my_flow_and_get_result_async(x: int) -> dict:
return await run_sync_in_worker_thread(run_my_flow_and_get_result, x)
Notes and options
- You can use any serialization you prefer (JSON, pickle, parquet). For large results, choose a format that matches your data. For tiny results, JSON is simplest.
- Keep cache_result_in_memory=False and persist_result=False at the flow/task level to prevent memory creep inside the subprocess.
- Error handling: check proc.exitcode. If you want richer error info, wrap your flow body in try/except and write an error file alongside the result.
- If you truly want to use Prefect’s result storage, you can set @flow(persist_result=True, result_storage="/path", result_serializer=...) but you’ll still need a way to communicate the storage key back to the parent. Passing an explicit output_path is usually simpler for client-side runs.
Docs for reference:
- Flow engine API (run_flow_in_subprocess): Flow engine API
- Async helpers: prefect.utilities.asyncutils
If you want, share your final result’s shape/size and I can suggest the best serializer and any concurrency tuning.R G
10/24/2025, 6:31 AMrun_flow_in_subprocess and directly use multiprocessing instead?Marvin
10/24/2025, 6:32 AMMarvin
10/24/2025, 6:32 AM# flows.py
from prefect import flow, task
@task(cache_result_in_memory=False, persist_result=False)
def step(x: int) -> dict:
# I/O-bound work here
return {"value": x}
@flow(cache_result_in_memory=False, persist_result=False, log_prints=False)
def my_flow(x: int) -> dict:
# ~20 tasks chained; only final result matters
return step(x)
# service.py
import multiprocessing as mp
from flows import my_flow
def _run_flow_worker(x: int, result_queue: "mp.Queue"):
try:
result = my_flow(x)
result_queue.put(("ok", result))
except Exception as e:
# send error info back to parent
result_queue.put(("err", repr(e)))
def run_my_flow_in_process(x: int):
ctx = mp.get_context("spawn") # cross-platform-safe
q = ctx.Queue()
p = ctx.Process(target=_run_flow_worker, args=(x, q))
p.start()
p.join()
if p.exitcode != 0:
# child crashed; no safe result to read
raise RuntimeError(f"Flow subprocess exited with code {p.exitcode}")
status, payload = q.get_nowait()
if status == "err":
raise RuntimeError(f"Flow failed: {payload}")
return payload
if __name__ == "__main__":
print(run_my_flow_in_process(123))
Using a file instead of a Queue (good for larger results)
# inside _run_flow_worker
import json, os, tempfile
...
result = my_flow(x)
with open(output_path, "w") as f:
json.dump(result, f)
result_queue.put(("ok", output_path))
Async service wrapper (don’t block the event loop)
from prefect.utilities.asyncutils import run_sync_in_worker_thread
async def run_my_flow_async(x: int):
return await run_sync_in_worker_thread(run_my_flow_in_process, x)
What you trade off vs run_flow_in_subprocess
- Pros of using multiprocessing directly:
- Full control over how you pass results (Queue, Pipe, files).
- No dependency on Prefect’s helper.
- Cons / things to handle yourself:
- You must set the spawn start method and main guard properly.
- You must manage result passing and error propagation.
- On Windows/macOS, anything not at top-level won’t be picklable for spawning.
Why run_flow_in_subprocess can be preferable
- It abstracts the cross-platform spawning details and passes parameters/context in a Prefect-friendly way. You still need an external channel (file/DB/object store) to get the final result, but it’s less boilerplate. API ref: Flow engine API
Given your case (I/O-bound tasks, only final result needed, memory creep concerns), either approach will give you the isolation you want. If you already have a pattern for exchanging results, multiprocessing directly is fine—just stick to spawn and keep results small/picklable or use files.R G
10/27/2025, 6:46 AMrun_sync_in_worker_thread.Marvin
10/27/2025, 6:46 AMMarvin
10/27/2025, 6:50 AMrun_sync_in_worker_thread and some practical examples.
What it does
- Offloads a blocking/synchronous function to a worker thread so you can call it from async code (e.g., async flows/tasks) without blocking the event loop.
- Uses a shared capacity limiter to prevent unbounded thread creation; you can override with your own limiter.
- Exceptions raised by the sync function propagate to the awaited caller.
- Cancellation of the awaiting task cancels the thread task as supported by AnyIO.
Import and signature
- Import path: from prefect.utilities.asyncutils import run_sync_in_worker_thread
- You must await it (use it in async contexts).
- Signature:
await run_sync_in_worker_thread(__fn, *args, **kwargs)
- `__fn`: the synchronous callable to run
- *args, `**kwargs`: passed to __fn
- Special keywords supported by the underlying AnyIO call include cancellable and limiter (see notes below).
Key notes
- Thread limiting: If you don’t provide a limiter, Prefect uses a shared limiter (via get_thread_limiter()) to avoid creating too many threads. You can pass your own AnyIO CapacityLimiter if you want to constrain a specific workload.
- Cancellation: If the parent async task is cancelled, AnyIO will cancel the worker thread task.
- Context: Prefect runtime context is available to the work executed in the thread.
- CPU-bound work: Threads won’t bypass the GIL. For CPU-bound workloads, prefer processes or specialized libraries; run_sync_in_worker_thread is best for blocking I/O.
- Name collisions: The underlying AnyIO function accepts cancellable and limiter as keywords. If your sync function also uses those keyword names, pass them positionally or via functools.partial to avoid conflicts.
Examples
1) Basic: call a blocking HTTP client (e.g., requests) from an async flow
from prefect import flow
from prefect.utilities.asyncutils import run_sync_in_worker_thread
import requests
def fetch(url, timeout=10):
return requests.get(url, timeout=timeout).json()
@flow
async def my_flow():
data = await run_sync_in_worker_thread(fetch, "<https://httpbin.org/get>", timeout=10)
return data
if __name__ == "__main__":
my_flow()
2) Use inside an async task
from prefect import task, flow
from prefect.utilities.asyncutils import run_sync_in_worker_thread
import time
def slow_blocking(n):
time.sleep(n)
return f"slept {n}s"
@task
async def do_work(n: int):
return await run_sync_in_worker_thread(slow_blocking, n)
@flow
async def my_flow():
res = await do_work(2)
print(res)
if __name__ == "__main__":
my_flow()
3) Limit concurrency with a custom limiter
from prefect import flow, task
from prefect.utilities.asyncutils import run_sync_in_worker_thread
from anyio import CapacityLimiter
import time
def io_heavy(x):
time.sleep(1)
return x
limiter = CapacityLimiter(5) # at most 5 threads for this workload
@task
async def work_item(x):
return await run_sync_in_worker_thread(io_heavy, x, limiter=limiter)
@flow
async def my_flow():
results = await work_item.map(range(20)) # fan-out but with a per-call limiter
return results
4) Avoid keyword collisions using functools.partial
from functools import partial
from prefect.utilities.asyncutils import run_sync_in_worker_thread
def fn_with_cancellable_kw(arg, cancellable=False):
# your function that happens to use 'cancellable'
...
partial_fn = partial(fn_with_cancellable_kw, cancellable=True)
result = await run_sync_in_worker_thread(partial_fn, "value")
Differences from related helpers
- `run_coro_as_sync`: Runs an async coroutine from sync code (the opposite direction). Use this when you are in sync code and need to call async code.Marvin
10/27/2025, 6:50 AMrun_sync_in_worker_thread for offloading sync work from existing async code.
Docs and references
- API reference: run_sync_in_worker_thread
- AnyIO threading docs (underlying implementation): AnyIO documentationR G
10/27/2025, 12:02 PMmap syntax here guarantee that all tasks in task1 are completed successfully before executing task2?Marvin
10/27/2025, 12:02 PMMarvin
10/27/2025, 12:09 PMtask1.map(...) (a list of futures) as an argument to task2.submit(...) makes Prefect resolve those futures before running task2. That means:
- Prefect waits for all mapped task1 runs to finish.
- If any of them failed, resolution raises and task2 will not start.
- Only if all upstream mapped runs complete successfully will task2 execute, receiving the resolved list of results.
If you want different behavior:
- Start task2 per item as soon as each corresponding task1 completes: map task2 as well
x1 = task1.map(x0)
x2 = task2.map(x1) # each task2[i] waits only for task1[i]
- Enforce ordering but not pass the results: use wait_for
x1 = task1.map(x0)
x2 = task2.submit(arg, wait_for=x1) # waits for all x1 to finish, doesn’t consume their values
- Proceed even if some upstreams failed: work with states instead of results
x1_states = task1.map(x0, return_state=True)
x2 = task2.submit(x1_states) # task2 can inspect which items succeeded/failed
Or, in your flow, resolve without raising:
x1 = task1.map(x0)
partial = x1.result(raise_on_failure=False) # list of results or exceptions
x2 = task2.submit(partial)
References
- Task API (submit/map): prefect.tasks
- Futures utilities: prefect.futures (see resolve_futures_to_results, resolve_futures_to_states, PrefectFutureList.result)R G
10/28/2025, 9:52 AMMarvin
10/28/2025, 9:52 AMMarvin
10/28/2025, 9:56 AMfrom prefect import task as prefect_task
from prefect import get_run_logger
from prefect.cache_policies import NO_CACHE
def task(*args, **kwargs):
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
kwargs.setdefault("log_prints", True) # 捕获 print
return prefect_task(*args, **kwargs)
@task() # 你的装饰器
def add(x, y):
print(f"[add] inputs: x={x}, y={y}")
result = x + y
print(f"[add] output: {result}")
return result
2) 用任务的状态 Hook 自动记录(通用、无需改动任务函数)
- 使用 on_completion Hook 可以拿到任务入参与结果,并统一写日志。
- Hook 的签名是 `(task, task_run, state)`;要获取返回值请在 Hook 里 `await state.result()`。
- 入参可通过 task_run.parameters 获取。
- 使用 get_run_logger() 打日志可进入 Prefect UI;你也可以用 `print`(配合 `log_prints=True`)。
参考文档:State change hooks
下面是将 Hook 集成到你自定义装饰器里的完整例子(自动打印输入与输出):
from prefect import task as prefect_task, get_run_logger
from prefect.cache_policies import NO_CACHE
# 完成后记录输入与输出
async def _log_on_completion(task, task_run, state):
logger = get_run_logger()
try:
<http://logger.info|logger.info>(f"[{task.name}] inputs: {task_run.parameters}")
except Exception as e:
logger.warning(f"[{task.name}] cannot read parameters: {e}")
try:
result = await state.result() # 获取任务返回值
<http://logger.info|logger.info>(f"[{task.name}] output: {result!r}")
except Exception as e:
logger.warning(f"[{task.name}] cannot resolve result: {e}")
# 可选:失败时也记录(包含异常)
async def _log_on_failure(task, task_run, state):
logger = get_run_logger()
try:
exc = await state.result(raise_on_failure=False)
logger.error(f"[{task.name}] failed: {exc!r}; inputs: {task_run.parameters}")
except Exception as e:
logger.error(f"[{task.name}] failed and could not read details: {e}")
def task(*args, **kwargs):
"""
你的自定义 task 装饰器:统一默认参数,并自动记录输入输出
"""
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
kwargs.setdefault("log_prints", True) # 如果你也想捕获 print
# 合并用户已传入的 hooks
on_completion = list(kwargs.pop("on_completion", []) or [])
on_completion.append(_log_on_completion)
kwargs["on_completion"] = on_completion
on_failure = list(kwargs.pop("on_failure", []) or [])
on_failure.append(_log_on_failure)
kwargs["on_failure"] = on_failure
return prefect_task(*args, **kwargs)
# 使用
@task()
def add(x, y):
return x + y
说明:
- 如果你更偏好“真的 print 到控制台”,确保 `log_prints=True`;这样在 Prefect 执行环境中也能看到这些信息同步到日志。
- 如果结果体积很大或不可序列化,建议在 Hook 里做截断或转字符串(如 repr(result) / `result!r`)。
- Hook 方式无需修改每个任务函数,更适合全局统一打印。
需要我帮你把这个封装成一个可复用的小模块,或根据你当前项目结构微调吗?R G
10/29/2025, 7:09 AMMarvin
10/29/2025, 7:09 AMMarvin
10/29/2025, 7:15 AM.submit() 返回的 PrefectFuture 在任务间传递、在 Python 分支里基于普通变量做 if/else、并把上游 Future 作为下游任务的入参,这些在 Prefect 3.x 都是推荐用法。Prefect 会在下游任务执行前自动把 Future 解析成实际值。
可以优化/注意的点与一个更“惯用”的写法如下:
- 若 aigc_video_res_list 是普通 Python 列表,分支没问题;如果它是某个任务产出的 Future,就需要先 .result() 再参与分支,或把“是否为空”的判断封装成一个小 task 返回布尔值再基于结果分支。
- 保证 preproc_video_effect 的输出和 blades_req 的输出在 schema 上对 preproc_merge 是一致的;否则在 else/if 两个分支里传给 preproc_merge 的数据结构会不一致。
- 更简洁的布尔表达式:`keep_ori_audio = bool(aigc_video_res_list)` 或 `keep_ori_audio = bool(len(...))`;不过同样前提是它不是 Future。
- 避免在 flow 里不必要地调用 `.result()`,会阻塞并降低并行度。你当前代码没有这么做,挺好。
建议的重写(等价但更紧凑):
video_effect_predata = preproc_video_effect.submit(
logid, planning_tts_json, tts_res, video_path_list_url, card_res_list, ttv_res_list, exp_id
)
# True 表示有 AIGC 结果,保留原音频;False 则需要调用 blades 做特效
keep_ori_audio = bool(aigc_video_res_list)
video_effect_res = (
video_effect_predata
if keep_ori_audio
else blades_req.submit(video_effect_predata) # 传 Future 没问题,Prefect 会自动解析
)
merge_predata = preproc_merge.submit(video_effect_res, voice_effect_res, keep_ori_audio)
merge_res = blades_req.submit(merge_predata)
如果存在以下情况,请相应调整:
- 若 aigc_video_res_list 是 Future:
keep_ori_audio = bool(aigc_video_res_list.result())
或者写一个小 task:
@task
def has_aigc(xs): return bool(xs)
keep_ori_audio = has_aigc.submit(aigc_video_res_list)
# 这里 keep_ori_audio 会是 Future,作为参数传给 preproc_merge 也会被自动解析
- 若 preproc_merge 对 video_effect_res 两种来源(preproc vs blades)有不同字段/格式,建议统一一个数据模型,或在进入 preproc_merge 之前做一次规整:
@task
def normalize_effect_data(data):
# 归一化到 preproc_merge 需要的 schema
return normalized
normalized_effect = normalize_effect_data.submit(video_effect_res)
merge_predata = preproc_merge.submit(normalized_effect, voice_effect_res, keep_ori_audio)
参考文档:
- 关于 Future 的自动解析与控制流分支:Task runners and futures
如果你能确认下:
- `aigc_video_res_list`、`voice_effect_res` 分别是普通值还是 Future?
- preproc_video_effect 与 blades_req 的输出在传给 preproc_merge 时是否完全一致?
我可以据此给出更精准的改写建议。R G
11/06/2025, 7:06 AMMarvin
11/06/2025, 7:06 AMMarvin
11/06/2025, 7:14 AMprefect CLI。
- 文档参考(建议随查随对照)
- Server 概念
- Server CLI
- 健康检查接口
- 日志与调试
- 自托管参考架构
一、快速健康检查(API 与就绪)
- 确认你客户端使用的 PREFECT_API_URL 指向的是 Server 的 /api 路径,例如 http://<host>:4200/api
- 健康/就绪探针
curl -i http://<host>:<port>/api/health
curl -i http://<host>:<port>/api/readiness
curl -sS http://<host>:<port>/api/admin/version
- 期望 /api/health 返回 200,`/api/readiness` 返回 200;能拿到 admin version 说明 API 正常连上并且 DB 基本可用
- 如果 behind Nginx/Ingress,确保反向代理转发了 `/api/*`,且没有把请求体/头限制得太严格(例如 Nginx 的 `client_max_body_size`)
二、打开服务端调试日志,重试并抓栈
- 临时用更高日志级别启动 Server,便于捕获 500 的栈:
PREFECT_LOGGING_ROOT_LEVEL=DEBUG \
PREFECT_SERVER_DATABASE_ECHO=true \
prefect server start --log-level DEBUG
然后在产生 500 的操作上重试一次(比如用你的客户端代码或 curl),观察终端输出的报错栈。容器/集群环境下使用:
- Docker: docker logs -f <prefect-server-container>
- Kubernetes: kubectl logs -f <pod> -c <container>
- 若你将 API 与服务进程分离运行,也可单独跑服务:
prefect server services start
并关注该进程的日志(调度、触发器、事件持久化等后台服务出错也会连带 500)
三、版本一致性检查(常见 500 根因)
- 新客户端对旧 Server 发请求时,若 API 发生不兼容变更,可能触发 500
# 分别在客户端和服务端环境执行
prefect version
若客户端版本较新,建议先升级服务端到同一大版本的最新补丁(或对齐版本)再试。
四、数据库检查与迁移
- 查看当前配置,确认 Server 连接到了你预期的数据库:
prefect config view --show-defaults | grep -i -E "DATABASE|API_URL|SERVER"
- 执行数据库迁移(很多 500 是因为表结构未升级/不完整)
prefect server database upgrade -y
如遇大型库迁移超时,可提高超时:
PREFECT_SERVER_DATABASE_TIMEOUT=30 prefect server database upgrade -y
- 使用 PostgreSQL 时,确保安装了 pg_trgm 扩展(Prefect 用于搜索/索引;缺失可能导致查询相关的 500)
CREATE EXTENSION IF NOT EXISTS pg_trgm;
如果是托管 PG,请在 DB 管理端确认是否支持该扩展。
- 若是 SQLite,确认路径可写、未被多个进程/容器并发写(并发写会报错)
五、后台服务状态(调度/触发器等)
- 列出并检查服务状态:
prefect server services list-services
- 确保未意外禁用关键服务(环境变量类似 PREFECT_SERVER_SERVICES_SCHEDULER_ENABLED=false 会禁用),或者你没有用 --no-services 启动 API 却忘了单独启动服务。
六、用最小化请求复现/区分 422 与 500
- 422 表示请求体校验失败(字段缺失/类型不对),500 才是服务端异常。你可以先做一个只读请求验证接口基本可用,例如列版本、健康检查(参考第一步)
- 如果你是通过 Python 客户端提交,建议加上客户端日志,或暂时用 curl/HTTP 客户端直接命中同一 REST 路径,对比响应状态码和响应体 detail
七、常见环境/网络问题
- 反向代理/负载均衡是否剥离了关键头(例如 `X-PREFECT-API-VERSION`)或限制了请求体、超时过短
- 服务器时间严重不准可能影响签名/调度(虽然罕见导致 500,但建议 NTP 同步)
- 若你水平扩展了多个 API 副本,确保它们共享同一个数据库,并且后台服务只跑一套(或采用推荐的架构)
你可以先按上面的顺序操作。为了更快定位,能否提供以下信息?
1) Server 与客户端的 prefect version 输出
2) 你的部署方式(本机/容器/Compose/K8s)、数据库类型(Postgres/SQLite)及数据库连接字符串隐藏敏感后的样子
3) PREFECT_API_URL 的实际值(隐藏域名/IP 也可,只要能判断是否带 `/api`)
4) 触发 500 时服务端日志的报错栈(开启 DEBUG 后的关键几行)
5) /api/health 与 /api/readiness 的返回结果/状态码
只要拿到上述信息,我们基本能快速锁定问题的根因。R G
11/06/2025, 8:08 AMMarvin
11/06/2025, 8:09 AMMarvin
11/06/2025, 8:09 AM# 先找容器名
docker ps --format 'table {{.Names}}\t{{.Image}}\t{{.Status}}\t{{.Ports}}'
# 跟随输出
docker logs -f <容器名>
# 只看最近 200 行并跟随
docker logs --tail=200 -f <容器名>
# 只看最近 30 分钟
docker logs --since=30m <容器名>
- Docker Compose 部署
# 查看所有服务日志
docker compose logs -f
# 只看某个服务(例如 api / services)
docker compose ps
docker compose logs -f <service_name>
在容器内查看(不推荐,但可以)
- 进入容器后,直接跟主进程(通常 PID 1)的 stdout/stderr
# stdout
tail -F /proc/1/fd/1
# stderr
tail -F /proc/1/fd/2
# 可用 ps 看下 PID 1 是什么
ps aux
说明:你 exec 进容器时看到的是一个新的 shell,会看不到历史日志;通过 /proc/1/fd/1/2 可以“窥视”主进程当前的输出。
提高日志详细程度(便于排查 500)
- 临时调高日志级别(修改 Compose 或启动命令后重启容器)
# docker run/compose 的 environment:
PREFECT_LOGGING_ROOT_LEVEL=DEBUG
PREFECT_SERVER_DATABASE_ECHO=true
# 或在 command 里
prefect server start --log-level DEBUG
- 日志配置参考:Prefect 日志配置
如果你不确定容器/服务名,发一下 docker ps 或 docker compose ps 的输出,我帮你定位应该看哪一个。R G
11/06/2025, 8:19 AMMarvin
11/06/2025, 8:19 AMMarvin
11/06/2025, 8:21 AMfrom prefect import flow, task, get_run_logger
from prefect.futures import as_completed
@task
def boom(i: int) -> int:
if i % 2:
raise ValueError(f"bad {i}")
return i * i
@flow
def main():
logger = get_run_logger()
futures = {i: boom.submit(i) for i in range(5)}
# 按提交顺序处理
for i, fut in futures.items():
try:
res = fut.result() # 这里失败会抛异常
<http://logger.info|logger.info>(f"boom({i}) -> {res}")
except Exception as exc:
logger.error(f"Task boom(i={i}) 失败, task_run_id={fut.task_run_id}, exc={exc}")
# 或者按完成顺序处理
for fut in as_completed(list(futures.values())):
try:
res = fut.result()
except Exception as exc:
# 通过反查 key 知道是哪个 i
i = next(k for k, v in futures.items() if v is fut)
logger.error(f"Task boom(i={i}) 失败, task_run_id={fut.task_run_id}, exc={exc}")
2) 使用 allow_failure,不让失败直接“炸掉”flow,转而通过 state 判断
from prefect import flow, task, get_run_logger, allow_failure
from prefect.states import get_state_exception
@task
def boom(i: int) -> int:
if i % 2:
raise ValueError(f"bad {i}")
return i * i
@flow
def main():
logger = get_run_logger()
futures = {i: allow_failure(boom.submit(i)) for i in range(5)}
for i, fut in futures.items():
st = fut.wait() # 拿到 State,不会抛异常
if st.is_failed():
exc = get_state_exception(st)
logger.error(f"boom({i}) 失败, task_run_id={fut.task_run_id}, exc={exc}")
else:
<http://logger.info|logger.info>(f"boom({i}) -> {st.result()}")
3) 需要更详细的 Task 运行信息(名称、标签等)
- 你可以用 task_run_id 去查询编排层的 TaskRun 对象
from prefect.client.orchestration import get_client
@flow
def main():
fut = boom.submit(1)
try:
fut.result()
except Exception:
with get_client() as client:
tr = client.read_task_run(fut.task_run_id)
# tr.name, tr.task_key, tr.tags, tr.parameters 等都可用
print(f"失败任务: name={tr.name}, id={tr.id}")
小结
- try/except 要包在 future.result() 周围;对“task 函数调用本身”做 try/except 是捕不到的,因为那一步只返回 future。
- 如果希望 flow 不因单个 task 失败而直接失败,用 `allow_failure()`,再基于 State 判断并提取异常。
- 用 future.task_run_id 可以唯一定位到失败的 TaskRun。R G
11/06/2025, 8:26 AMMarvin
11/06/2025, 8:26 AMMarvin
11/06/2025, 8:28 AMFailed 这个对象当作“正常返回值”,不会把 task 状态标记为失败,更不能让 flow “自动知道哪个 task 失败了”。如果要让 task 真的失败,请在 task 内部直接 raise 异常。
- 想在不逐个 future.result() 的情况下“批量”拿到失败任务,可以在提交时统一收集 future,最后一次性解析所有 state 来筛选失败项;或者让任务自己在失败时记录名字。
可选做法一:统一收集 + 批量解析状态(推荐)
from prefect import flow, task, get_run_logger, allow_failure
from prefect.futures import resolve_futures_to_states
from prefect.client.orchestration import get_client
@task
def boom(i: int):
if i % 2:
raise ValueError(f"bad {i}")
return i*i
@flow
def main(n: int = 1000):
logger = get_run_logger()
# 1) 统一收集 future,并允许失败不阻塞后续(关键)
futures = [allow_failure(boom.submit(i)) for i in range(n)]
# 2) 一次性把所有 future 解析为 State
states = resolve_futures_to_states(futures)
# 3) 找到失败的 futures(不需要逐个 result())
failed = [f for f, st in states.items() if st.is_failed()]
<http://logger.info|logger.info>(f"共失败 {len(failed)} 个 task")
# 4) 想拿“任务名”,可用 task_run_id 反查
if failed:
from prefect.states import get_state_exception
with get_client() as client:
for f in failed:
tr = client.read_task_run(f.task_run_id)
exc = get_state_exception(states[f])
logger.error(f"Task 失败: name={tr.name}, id={tr.id}, exc={exc}")
要点:
- 用 allow_failure(...) 包装提交,保证即使 task 失败,flow 也能跑完并在末尾做集中统计。
- resolve_futures_to_states 可一次性得到所有 State,再筛 is_failed() 即可。
- 通过 future.task_run_id → client.read_task_run(...) 能拿到任务名等元数据。
可选做法二:任务内自报姓名并失败(无需集中扫描)
from prefect import task, get_run_logger
from prefect.context import get_run_context
@task
def wrapped(...):
logger = get_run_logger()
ctx = get_run_context() # 在 task 中可拿到 TaskRun 上下文
try:
return real_work(...)
except Exception as exc:
logger.error(f"Task 失败: name={ctx.task_run.name}, id={ctx.task_run.id}, exc={exc}")
raise # 仍然让任务真正失败
这个方式会在每个失败发生时立即记录任务名,flow 最终也会因为失败而失败;不需要在 flow 里扫描 futures。
可选做法三:按完成顺序处理(仍需遍历,但不用“逐个写 try/except”)
from prefect.futures import as_completed
futures = [boom.submit(i) for i in range(n)]
for fut in as_completed(futures):
try:
_ = fut.result()
except Exception as exc:
# 通过 fut.task_run_id 反查名字
...
补充说明
- 再次强调:返回 prefect.states.Failed 不会改变任务运行状态,想失败必须 `raise`。
- 如果你的需求是“flow 不因个别任务失败而中断,但最终能一次性拿到所有失败任务的名字”,做法一最贴合:`allow_failure` + 批量解析状态 + 用 task_run_id 反查名字。
如果你方便的话,告诉我你的任务启动方式(for 循环、批量 submit、递归等),我可以给你贴合你代码结构的最小改造示例。R G
11/10/2025, 6:54 AMMarvin
11/10/2025, 6:54 AMMarvin
11/10/2025, 7:05 AMP1D 或 `P3D`(ISO8601 时长,默认是 `P7D`)
- 这对数据库增长影响巨大,因为 events 表是增长最快的
- 文档:Settings reference
2) 清理已有的历史数据(释放磁盘空间)
- 官方没有内置“一键批量清理”CLI,需要用数据库维护或脚本化清理
- 推荐顺序:events → logs → flow/task runs → artifacts
- 维护指南:Database maintenance
A. 先找出最大表(诊断)
你可以先在 PostgreSQL 中确认体量最大的表,再决定清理策略:
-- 看前 20 大表(含索引)
SELECT relname AS table_name,
pg_size_pretty(pg_total_relation_size(oid)) AS total_size
FROM pg_class
WHERE relkind = 'r'
ORDER BY pg_total_relation_size(oid) DESC
LIMIT 20;
-- 每张表的行数估计
SELECT relname, n_live_tup
FROM pg_stat_user_tables
ORDER BY n_live_tup DESC
LIMIT 20;
在 Prefect 的生产经验中,通常是 event > log > `flow_run`/`task_run` > _state 系列表。
B. events 清理
- 缩短 PREFECT_SERVER_EVENTS_RETENTION_PERIOD 后,服务端的 Event Persister 会按保留期自动删除过期事件;但对已积累的超大数据,需要配合一次性批量清理
- 如果你能接受直接在 DB 删除,示例(请先用 SELECT 验证范围再执行 DELETE):
-- 预览将要删除的事件范围(举例:删除 7 天前)
SELECT COUNT(*) FROM event WHERE occurred < now() - interval '7 days';
-- 真正删除
DELETE FROM event WHERE occurred < now() - interval '7 days';
注意:表名、时间列名请以你的数据库实际为准(Prefect ORM 模型名是 Event,通常表名为 event,时间列常为 occurred)。先用上面的“表体积/列名检查”确认。
C. logs 清理
- 若之前将日志上报到 API,logs 表会很大
- Prefect server 层提供按过滤器删除日志的能力(Python API,服务内调用),也可直接 DB 清理:
-- 预览
SELECT COUNT(*) FROM log WHERE "timestamp" < now() - interval '7 days';
-- 删除
DELETE FROM log WHERE "timestamp" < now() - interval '7 days';
同样,注意表/列名以你的 DB 实际结构为准;timestamp 列通常命名为 timestamp。
D. flow/task runs 和 states 清理
- 你可以按开始时间/结束时间删除很久以前的 flow run;这通常会通过外键级联或触发器清理部分关联记录,但不要假设所有依赖(比如日志、artifacts)都会自动清
- 更稳妥的方式是分层分表清:
1) 先删 events、logs
2) 再删很老的 flow run(以及 task run 和 state)
- 示例(仅供思路,先用 SELECT 验证,列名以实际为准):
-- 预览 30 天前完成的 flow run
SELECT COUNT(*) FROM flow_run
WHERE end_time < now() - interval '30 days';
-- 删除
DELETE FROM flow_run
WHERE end_time < now() - interval '30 days';
如果你更倾向于走 Prefect 的 Python/REST API,可以在维护脚本中用 PrefectClient 列出很早的 flow run,然后逐个 DELETE,但吞吐会比直接 SQL 慢很多。Python Server 模型也提供了删除函数(用于服务内部),例如:
- logs:`prefect.server.models.logs.delete_logs(log_filter=...)`
- flow run:`prefect.server.models.flow_runs.delete_flow_run(flow_run_id=...)`
- artifacts:`prefect.server.models.artifacts.delete_artifact(artifact_id=...)`
API 参考:
- 日志与过滤器:prefect.server.models.logs
- Flow runs:prefect.server.models.flow_runs
- Artifacts:prefect.server.models.artifacts
- REST 概览:DELETE /flow-runs/{id}
E. artifacts 清理(如不需要保留)
- 如果你创建了大量 artifacts(度量、表格、结果预览等),可以先确认是否需要长期保留
- 没有内置批量保留策略,需要脚本遍历并删除(按时间或按 key 前缀等)
F. 删除后的 PostgreSQL 维护
- 执行 `VACUUM (ANALYZE)`,必要时 VACUUM FULL 或使用 `pg_repack`(推荐生产用 `pg_repack`,减少锁)
- 检查膨胀的索引,必要时重建索引
- 文档:Database maintenance
3) 长期优化建议
- 设定严格的保留窗口
- 对于每天 50 万 flow 的量级,建议:
- events:P1D ~ P3D
- logs:尽可能关闭 API 日志或仅保留 P1D ~ P3DMarvin
11/10/2025, 7:05 AMenvironment:
# 强烈建议:减少或关闭日志上报
PREFECT_LOGGING_TO_API_ENABLED: "false"
PREFECT_LOGGING_LEVEL: "WARNING"
# 缩短 events 保留期
PREFECT_SERVER_EVENTS_RETENTION_PERIOD: "P1D"
需要你补充的信息(好帮你定更精细的清理脚本)
- 你当前的 Prefect 版本和部署方式(容器/裸机)
- 是否可以关闭 API 日志?如果不行,最短可接受的日志保留期是多久?
- 对 flow run 历史记录的合规要求(需要保留几天/几个月?)
- 目前 events 保留期配置是多少(默认 7 天吗)?
- 是否允许直接对 PostgreSQL 执行 SQL 清理,还是只能通过 API/脚本?
参考文档
- Settings 参考(包含日志与 events 保留设置):Settings reference
- 日志定制:Logging customization
- 数据库维护与清理:Database maintenance
- Prefect Server 概念:Server concepts
如果你愿意,我可以根据你的保留需求给你一份可执行的“每日清理 SQL/脚本”,并评估一次性清理会释放多少空间。R G
11/10/2025, 11:07 AMMarvin
11/10/2025, 11:07 AMMarvin
11/10/2025, 11:08 AM2d 改成 P2D 才会生效。此外,最好用服务端键名,并确保是在运行 Prefect Server 的同一环境里设置并重启服务。
建议用法
- 正确设置(建议使用服务端键名):
prefect config set PREFECT_SERVER_EVENTS_RETENTION_PERIOD="P2D"
- 你用的键名 PREFECT_EVENTS_RETENTION_PERIOD 在 3.x 中有别名支持,一般也能识别,但推荐显式使用 `PREFECT_SERVER_EVENTS_RETENTION_PERIOD`。
- 设置后执行校验,确认没有解析错误:
prefect config validate
prefect config view | grep EVENTS_RETENTION_PERIOD
- 如果是用 Docker/容器跑的 Prefect Server,需要把这个环境变量加到 server 容器的 environment 中,然后重启容器;在本机跑 prefect config set 不会影响容器内的进程。
- 修改只影响 events 表的保留期(默认 7 天);对 logs、flow/task runs 不起作用。要进一步控盘请同时减少或关闭日志上报(`PREFECT_LOGGING_TO_API_ENABLED=false` 或提高 `PREFECT_LOGGING_LEVEL=WARNING/ERROR`),并对历史 logs/flow runs 做一次性清理和 VACUUM。
文档
- 设置参考(含 events 保留期与日志设置):Settings reference
- 日志定制与上报:Logging customization
- 数据库维护与清理:Database maintenance
需要的话我可以根据你的部署方式(裸机或 docker-compose)给出具体的环境变量配置和重启步骤。R G
11/10/2025, 11:54 AMMarvin
11/10/2025, 11:54 AMMarvin
11/10/2025, 11:56 AMPREFECT_SERVER_EVENTS_RETENTION_PERIOD="P2D" 或更短
- 减少/关闭日志上报到 API(强烈建议):
- `PREFECT_LOGGING_TO_API_ENABLED=false`,或至少 PREFECT_LOGGING_LEVEL=WARNING
- 设置后重启 Prefect Server 容器/进程
- 文档:
- 设置参考:Settings reference
- 日志定制:Logging customization
二、一次性清理:按表的体量从大到小、由依赖到主表安全删除
注意事项
- 大表删除务必“分批”执行,避免长事务和表级锁
- 每一步先用 SELECT 预览数量,再 DELETE
- 删除后执行 VACUUM / REINDEX,或考虑 pg_repack 以更少锁重整表
- 下述列名/表名是 Prefect 3 的常见命名,请先用信息架构查询确认
0) 辅助:找出外键关系(用于后续 orphan 清理)
-- 哪些表引用了某个表(例如 event_resources):
SELECT
tc.table_schema, tc.table_name, kcu.column_name,
ccu.table_name AS foreign_table_name, ccu.column_name AS foreign_column_name,
tc.constraint_name
FROM information_schema.table_constraints AS tc
JOIN information_schema.key_column_usage AS kcu
ON tc.constraint_name = kcu.constraint_name
JOIN information_schema.constraint_column_usage AS ccu
ON ccu.constraint_name = tc.constraint_name
WHERE tc.constraint_type = 'FOREIGN KEY'
AND ccu.table_name = 'event_resources';
1) events(最大头)和 event_resources
- 先批量删除老的 events(按 occurred 时间);Prefect 的事件保留服务会按保留期自动删,但历史积压需要一次性清
预览与分批删除(示例:删除 3 天前)
-- 预览
SELECT COUNT(*) FROM event WHERE occurred < now() - interval '3 days';
-- 分批删除(每批 100 万行,循环执行多次)
DELETE FROM event
WHERE occurred < now() - interval '3 days'
ORDER BY occurred
LIMIT 1000000;
- event_resources 清理:如果 schema 中 resources 通过关联关系被 events 引用,删完 events 后会有孤儿资源。可做 orphan 清理(根据你的外键结构调整联结):
-- 示例:如果 event 表通过 event.resource_id 指向 event_resources(id)
DELETE FROM event_resources er
WHERE NOT EXISTS (
SELECT 1 FROM event e WHERE e.resource_id = er.id
);
如果是多对多关系(中间表),就改成“NOT EXISTS 关联中间表”的写法。请用“步骤0”的外键查询确定真实关系。
2) logs
- 如果已经决定关闭或缩短日志保留,这里一次性清理历史日志
-- 预览
SELECT COUNT(*) FROM log WHERE "timestamp" < now() - interval '7 days';
-- 分批删除
DELETE FROM log
WHERE "timestamp" < now() - interval '7 days'
ORDER BY "timestamp"
LIMIT 100000;
3) flow/task runs 与 states(你的 task_run_state 和 task_run 很大)
最佳做法是基于“旧 flow run”作为边界来清理依赖对象,因为它能最干净地刮掉一整个执行链条的数据。
- 建议选一个业务可接受的保留期(例如只保留近 7~30 天的已结束 flow run)
- 按顺序删除:task_run_state → task_run → flow_run_state → flow_run(如果外键没有 ON DELETE CASCADE,按这个顺序最安全;若存在 CASCADE,可直接删 flow_run)
示例(以 end_time 早于 30 天为准;若有未结束的运行请避开):
-- 先找出要删除的旧 flow_run
WITH old_flows AS (
SELECT id FROM flow_run
WHERE end_time IS NOT NULL
AND end_time < now() - interval '30 days'
LIMIT 50000 -- 分批
)
-- 预览计数
SELECT COUNT(*) FROM old_flows;
-- 先删与旧 flow 相关的日志(如果没在前面按时间统一清过)
DELETE FROM log
WHERE flow_run_id IN (SELECT id FROM old_flows);
-- 删 task_run_state(依赖 task_run)
DELETE FROM task_run_state
WHERE task_run_id IN (
SELECT id FROM task_run
WHERE flow_run_id IN (SELECT id FROM old_flows)
);
-- 删 task_run
DELETE FROM task_run
WHERE flow_run_id IN (SELECT id FROM old_flows);
-- 删 flow_run_state(依赖 flow_run)
DELETE FROM flow_run_state
WHERE flow_run_id IN (SELECT id FROM old_flows);
-- 删 flow_run
DELETE FROM flow_run
WHERE id IN (SELECT id FROM old_flows);
如果你的系统对“只保留最近一次状态”有诉求(即不想删整个老 flow,而是清理老状态),可以改成“对每个 task_run 仅保留最新一条 task_run_state”。但考虑到你数据规模很大,这类去重型清理复杂且慢,通常不如“按时间删整批旧 flow runs”来得有效和稳妥。
4) 执行维护(释放空间)
- 对大表执行:
VACUUM (ANALYZE) event;
VACUUM (ANALYZE) event_resources;
VACUUM (ANALYZE) log;
VACUUM (ANALYZE) task_run_state;
VACUUM (ANALYZE) task_run;
VACUUM (ANALYZE) flow_run_state;
VACUUM (ANALYZE) flow_run;
- 如果需要立即回收物理空间、且能接受更强锁:`VACUUM FULL <table>;`
- 生产上更推荐用 pg_repack 对热点大表做在线重整(更少锁)
- 维护指南:Database maintenance
三、长期策略与建议
- 保留期建议(你的规模:50 万 flow/天)
- events:1~3 天
- logs:尽量关闭 API 上报;若必须保留,1~3 天Marvin
11/10/2025, 11:56 AMR G
11/10/2025, 12:28 PMMarvin
11/10/2025, 12:28 PMMarvin
11/10/2025, 12:34 PMimport os
import time
from datetime import datetime, timedelta, timezone
from prefect.client.orchestration import SyncPrefectClient
from prefect.client.schemas.filters import (
FlowRunFilter,
FlowRunFilterStartTime,
FlowRunFilterState,
FlowRunFilterStateType,
ArtifactFilter,
ArtifactFilterFlowRunId,
)
from prefect.client.schemas.objects import StateType
# 配置
# 请确保运行环境已设置 PREFECT_API_URL 指向你的自托管 Server,例如:
# os.environ["PREFECT_API_URL"] = "http://<your-server-host>:4200/api"
RETENTION_DAYS_FOR_RUNS = 30 # 仅保留近 30 天的 flow runs(示例)
BATCH_SIZE = 500 # 每批处理条数(按环境调节)
SLEEP_BETWEEN_BATCHES_SEC = 0.2 # 批次间隔,避免打爆 API/DB
DRY_RUN = True # 先 dry-run 看看命中数量,再设为 False 执行删除
# 仅删除“已结束”的 runs,避免误删进行中的 runs
END_STATES = [
StateType.COMPLETED,
StateType.FAILED,
StateType.CANCELLED,
StateType.CRASHED,
]
def delete_artifacts_for_run(client: SyncPrefectClient, flow_run_id):
# 按 flow_run_id 过滤 artifacts 并删除
# 分页拉取,直到为空
total_deleted = 0
while True:
arts = client.read_artifacts(
artifact_filter=ArtifactFilter(
flow_run_id=ArtifactFilterFlowRunId.any_([flow_run_id])
),
limit=BATCH_SIZE,
offset=0, # 我们每次都从 0 拉,因为会不断删除,集合会缩小
)
if not arts:
break
for a in arts:
if DRY_RUN:
print(f"[DRY] would delete artifact {a.id} (key={getattr(a, 'key', None)}) for flow_run {flow_run_id}")
else:
client.delete_artifact(a.id)
total_deleted += 1
if SLEEP_BETWEEN_BATCHES_SEC:
time.sleep(SLEEP_BETWEEN_BATCHES_SEC)
return total_deleted
def main():
cutoff = datetime.now(timezone.utc) - timedelta(days=RETENTION_DAYS_FOR_RUNS)
with SyncPrefectClient() as client:
print(f"Start cleanup. Retention for runs: {RETENTION_DAYS_FOR_RUNS} days; cutoff={cutoff.isoformat()}; DRY_RUN={DRY_RUN}")
total_runs_deleted = 0
total_artifacts_deleted = 0
while True:
# 每次取一批“早于 cutoff 且已结束”的 flow runs
runs = client.read_flow_runs(
flow_run_filter=FlowRunFilter(
start_time=FlowRunFilterStartTime(before_=cutoff),
state=FlowRunFilterState(
type=FlowRunFilterStateType.any_(END_STATES)
),
),
limit=BATCH_SIZE,
offset=0, # 结合删除操作,始终从 0 开始拿
)
if not runs:
break
print(f"Found {len(runs)} old flow runs to delete in this batch...")
for r in runs:
# 先删 artifacts(如果你从不使用 artifacts,也可以跳过这步)
deleted_art_count = delete_artifacts_for_run(client, r.id)
total_artifacts_deleted += deleted_art_count
# 再删 flow run(一般会级联删除该 run 下的 task runs 与 states)
if DRY_RUN:
print(f"[DRY] would delete flow run {r.id} (name={getattr(r, 'name', None)}, start={getattr(r, 'start_time', None)})")
else:
client.delete_flow_run(r.id)
total_runs_deleted += 1
if SLEEP_BETWEEN_BATCHES_SEC:
time.sleep(SLEEP_BETWEEN_BATCHES_SEC)
print(f"Cleanup done. runs_deleted={total_runs_deleted}, artifacts_deleted={total_artifacts_deleted}, DRY_RUN={DRY_RUN}")
if __name__ == "__main__":
main()
要点和注意Marvin
11/10/2025, 12:34 PMPREFECT_LOGGING_TO_API_ENABLED=false 或 PREFECT_LOGGING_LEVEL=WARNING/ERROR
- 文档:Logging customization
- 级联删除说明:删除 flow run 通常会清理其 task runs 和 states(由服务端实现级联/事务)。但不同版本/迁移状态下行为可能有差异;如果你观察到残留,可在脚本里增加“按 flow_run_id 读取 task runs 并逐个删除”的兜底逻辑。
- 大量删除后,请在数据库侧执行 VACUUM/ANALYZE 或使用 pg_repack 以释放空间并更新统计信息:
- Database maintenance
文档参考
- 客户端用法(async/sync):Advanced: API client
- 过滤器(FlowRun/Artifact/Log 等):Client filter schemas
- Flow run 删除 REST 端点:DELETE /flow-runs/{id}
- Artifact 删除 REST 端点:DELETE /artifacts/{id}
如果你希望改为异步版本(PrefectClient)或加入“仅删除最近 N 条 state、保留最新一条 state”的策略,我可以基于你的保留需求再给你一份增强脚本。R G
11/10/2025, 12:44 PMMarvin
11/10/2025, 12:44 PMMarvin
11/10/2025, 12:45 PM-- 建议先确保有时间列索引,加速筛选
-- CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_ts ON log ("timestamp");
WITH to_delete AS (
SELECT id
FROM log
WHERE "timestamp" < now() - interval '30 days'
ORDER BY "timestamp"
LIMIT 100000
)
DELETE FROM log l
USING to_delete d
WHERE l.id = d.id;
如果没有合适的主键,可以用 ctid(单语句内唯一):
WITH to_delete AS (
SELECT ctid
FROM log
WHERE "timestamp" < now() - interval '30 days'
ORDER BY "timestamp"
LIMIT 100000
)
DELETE FROM log l
USING to_delete d
WHERE l.ctid = d.ctid;
循环批删(一次 10 万行,直到删不动为止)
DO $$
DECLARE
rows_deleted integer;
BEGIN
LOOP
WITH to_delete AS (
SELECT ctid
FROM log
WHERE "timestamp" < now() - interval '30 days'
LIMIT 100000
)
DELETE FROM log l
USING to_delete d
WHERE l.ctid = d.ctid;
GET DIAGNOSTICS rows_deleted = ROW_COUNT;
EXIT WHEN rows_deleted = 0;
PERFORM pg_sleep(0.1); -- 批次间隙,避免长时间锁
END LOOP;
END
$$;
同理,之前给的 events、event_resources、task_run_state、task_run 的分批删除也要用这种 CTE/子查询写法。例如 events(按 occurred 时间):
WITH to_delete AS (
SELECT ctid
FROM event
WHERE occurred < now() - interval '3 days'
LIMIT 100000
)
DELETE FROM event e
USING to_delete d
WHERE e.ctid = d.ctid;
建议
- 大批量删除前先创建时间列索引(CONCURRENTLY,避免长锁),如 `"timestamp"`、`occurred`。
- 每次删除后跑 `VACUUM (ANALYZE)`,需要立即回收空间可用 pg_repack 或在维护窗口用 `VACUUM FULL`。
- Prefect 官方维护指引:Database maintenance
如果你希望我把这些 SQL 改造成针对你所有大表(events、event_resources、task_run_state、task_run、log)的完整“批量清理脚本”,告诉我各表的保留期,我直接给你一份可执行版本。R G
11/10/2025, 12:47 PMMarvin
11/10/2025, 12:47 PMMarvin
11/10/2025, 12:48 PMPREFECT_LOGGING_LEVEL=WARNING
一、用 Prefect 3 Client 删除“早于 15 天且已结束”的 flow runs(附带删 artifacts)
说明:
- 仅清理早于 15 天、且处于终态(COMPLETED/FAILED/CANCELLED/CRASHED)的 flow runs
- 先删 artifacts,再删 flow run(通常会一并清理 task runs 和 states)
- 这是同步客户端版本,适合 cron/运维作业;先 DRY_RUN 验证
import os
import time
from datetime import datetime, timedelta, timezone
from prefect.client.orchestration import SyncPrefectClient
from prefect.client.schemas.filters import (
FlowRunFilter,
FlowRunFilterStartTime,
FlowRunFilterState,
FlowRunFilterStateType,
ArtifactFilter,
ArtifactFilterFlowRunId,
)
from prefect.client.schemas.objects import StateType
# 请确保运行环境有 PREFECT_API_URL 指向你的自托管 Server,如:
# os.environ["PREFECT_API_URL"] = "http://<server-host>:4200/api"
RETENTION_DAYS_FOR_RUNS = 15
BATCH_SIZE = 500
SLEEP_BETWEEN_BATCHES_SEC = 0.2
DRY_RUN = True # 先 dry-run;确认无误后改为 False 执行删除
END_STATES = [
StateType.COMPLETED,
StateType.FAILED,
StateType.CANCELLED,
StateType.CRASHED,
]
def delete_artifacts_for_run(client: SyncPrefectClient, flow_run_id):
total_deleted = 0
while True:
arts = client.read_artifacts(
artifact_filter=ArtifactFilter(
flow_run_id=ArtifactFilterFlowRunId.any_([flow_run_id])
),
limit=BATCH_SIZE,
offset=0,
)
if not arts:
break
for a in arts:
if DRY_RUN:
print(f"[DRY] would delete artifact {a.id} for flow_run {flow_run_id}")
else:
client.delete_artifact(a.id)
total_deleted += 1
if SLEEP_BETWEEN_BATCHES_SEC:
time.sleep(SLEEP_BETWEEN_BATCHES_SEC)
return total_deleted
def main():
cutoff = datetime.now(timezone.utc) - timedelta(days=RETENTION_DAYS_FOR_RUNS)
with SyncPrefectClient() as client:
print(f"Start cleanup: runs<{RETENTION_DAYS_FOR_RUNS}d; cutoff={cutoff.isoformat()}; DRY_RUN={DRY_RUN}")
total_runs_deleted = 0
total_artifacts_deleted = 0
while True:
runs = client.read_flow_runs(
flow_run_filter=FlowRunFilter(
start_time=FlowRunFilterStartTime(before_=cutoff),
state=FlowRunFilterState(
type=FlowRunFilterStateType.any_(END_STATES)
),
),
limit=BATCH_SIZE,
offset=0,
)
if not runs:
break
for r in runs:
total_artifacts_deleted += delete_artifacts_for_run(client, r.id)
if DRY_RUN:
print(f"[DRY] would delete flow run {r.id} (start={getattr(r, 'start_time', None)})")
else:
client.delete_flow_run(r.id)
total_runs_deleted += 1
if SLEEP_BETWEEN_BATCHES_SEC:
time.sleep(SLEEP_BETWEEN_BATCHES_SEC)
print(f"Done. runs_deleted={total_runs_deleted}, artifacts_deleted={total_artifacts_deleted}, DRY_RUN={DRY_RUN}")
if __name__ == "__main__":
main()
参考:
- API Client:Advanced: API client
- 过滤器:Client filter schemas
- Flow run 删除端点:DELETE /flow-runs/{id}
- Artifact 删除端点:DELETE /artifacts/{id}
二、用 SQL 清理 events、event_resources、logs(15 天)
说明:
- PostgreSQL 的 DELETE 不能直接 ORDER BY/LIMIT,需用 CTE 挑一批主键/ctid 再删Marvin
11/10/2025, 12:48 PM-- events 时间列
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_event_occurred ON event (occurred);
-- logs 时间列
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_ts ON log ("timestamp");
-- flow/task 依赖键
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_task_run_flow ON task_run (flow_run_id);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_task_state_task ON task_run_state (task_run_id);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_flow_run_end ON flow_run (end_time);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_flow_state_flow ON flow_run_state (flow_run_id);
2) events(删除 15 天前)
-- 重复执行,直到 0 行被删
WITH to_delete AS (
SELECT ctid
FROM event
WHERE occurred < now() - interval '15 days'
LIMIT 100000
)
DELETE FROM event e
USING to_delete d
WHERE e.ctid = d.ctid;
3) event_resources(删除孤儿资源)
注意:具体列名需以你的 schema 为准(常见为 event.resource_id → event_resources.id)。如果你的事件资源是多对多中间表,请据实调整 NOT EXISTS 的联结。
-- 删掉不再被任何 event 引用的资源行
DELETE FROM event_resources er
WHERE NOT EXISTS (
SELECT 1 FROM event e WHERE e.resource_id = er.id
);
4) logs(删除 15 天前)
WITH to_delete AS (
SELECT ctid
FROM log
WHERE "timestamp" < now() - interval '15 days'
LIMIT 100000
)
DELETE FROM log l
USING to_delete d
WHERE l.ctid = d.ctid;
5) flow/task(可选:如果你还想统一把早于 15 天的旧运行整批清掉)
说明:以 flow_run.end_time 为依据,只处理已结束的运行;按照“依赖→主记录”的顺序删除,避免外键约束失败。
-- 先挑出旧的 flow_run(分批)
WITH old_flows AS (
SELECT id
FROM flow_run
WHERE end_time IS NOT NULL
AND end_time < now() - interval '15 days'
LIMIT 50000
)
-- 如还没清 logs,可先清这些 runs 的日志
DELETE FROM log
WHERE flow_run_id IN (SELECT id FROM old_flows);
-- 先删 task_run_state
DELETE FROM task_run_state
WHERE task_run_id IN (
SELECT id FROM task_run
WHERE flow_run_id IN (SELECT id FROM old_flows)
);
-- 再删 task_run
DELETE FROM task_run
WHERE flow_run_id IN (SELECT id FROM old_flows);
-- 再删 flow_run_state
DELETE FROM flow_run_state
WHERE flow_run_id IN (SELECT id FROM old_flows);
-- 最后删 flow_run
DELETE FROM flow_run
WHERE id IN (SELECT id FROM old_flows);
6) 维护(释放空间)
VACUUM (ANALYZE) event;
VACUUM (ANALYZE) event_resources;
VACUUM (ANALYZE) log;
VACUUM (ANALYZE) task_run_state;
VACUUM (ANALYZE) task_run;
VACUUM (ANALYZE) flow_run_state;
VACUUM (ANALYZE) flow_run;
如需立即回收磁盘、且能接受更强锁,可用 `VACUUM FULL`;生产更推荐 pg_repack 等方案。
三、别忘了“止血”配置(避免再次长胖)
- 事件保留期:将运行 Prefect Server 的环境设为
PREFECT_SERVER_EVENTS_RETENTION_PERIOD="P15D"
- 日志上报:尽量关闭或降级
PREFECT_LOGGING_TO_API_ENABLED=false
# 或
PREFECT_LOGGING_LEVEL=WARNING
参考:
- 设置参考:Settings reference
- 日志定制:Logging customization
需要我帮你把上述 SQL 打包成可重复执行的脚本(例如 psql 可执行 + 循环批删),或者把 Python Client 版换成异步并增加并行度的话,告诉我你偏好的运行方式(裸机/容器、是否可安装 pg_repack、可接受的批次大小/限流)。我可以按你的环境生成一份一键脚本。R G
11/10/2025, 1:03 PMMarvin
11/10/2025, 1:03 PMMarvin
11/10/2025, 1:04 PM-- flow runs 及其关联(通常会级联 task runs / states)
VACUUM (ANALYZE) flow_run;
VACUUM (ANALYZE) flow_run_state;
VACUUM (ANALYZE) task_run;
VACUUM (ANALYZE) task_run_state;
-- 如果你也清了 logs / events
VACUUM (ANALYZE) log;
VACUUM (ANALYZE) event;
VACUUM (ANALYZE) event_resources;
需要立即缩小文件体积(谨慎)
-- 锁表强、影响写入;更推荐用 pg_repack 做在线重整
VACUUM FULL flow_run;
VACUUM FULL task_run;
VACUUM FULL task_run_state;
VACUUM FULL flow_run_state;
补充
- Prefect Client 不会帮你做 VACUUM;这一步需要在 PostgreSQL 上单独执行。
- 长期建议:建立定期清理+VACUUM/pg_repack 的维护作业,避免再次快速膨胀。
- 相关参考:Database maintenanceR G
11/10/2025, 1:05 PMMarvin
11/10/2025, 1:05 PMMarvin
11/10/2025, 1:06 PMVACUUM (ANALYZE) flow_run;
VACUUM (ANALYZE) flow_run_state;
VACUUM (ANALYZE) task_run;
VACUUM (ANALYZE) task_run_state;
VACUUM (ANALYZE) log;
VACUUM (ANALYZE) event;
VACUUM (ANALYZE) event_resources;
- 需要立刻降磁盘:用 pg_repack 或在窗口期做 `VACUUM FULL <table>`。
2) 让 autovacuum 更“积极”,特别是 events/log 这类高写表
- 表级参数(不重启即可生效):
ALTER TABLE event
SET (autovacuum_vacuum_scale_factor = 0.01,
autovacuum_vacuum_threshold = 10000,
autovacuum_analyze_scale_factor = 0.02);
ALTER TABLE log
SET (autovacuum_vacuum_scale_factor = 0.01,
autovacuum_vacuum_threshold = 10000,
autovacuum_analyze_scale_factor = 0.02);
- 全局参数(需改 postgresql.conf 并重启或 reload)示例:
autovacuum = on
autovacuum_max_workers = 6
autovacuum_naptime = '10s'
autovacuum_vacuum_cost_limit = 2000
- 开启日志便于观察是否触发及时:
ALTER SYSTEM SET log_autovacuum_min_duration = '0';
SELECT pg_reload_conf();
3) 监控是否“跟得上”
-- 关注死元组和最近的自动/手动 VACUUM 时间
SELECT relname, n_dead_tup, last_autovacuum, last_vacuum, autovacuum_count, vacuum_count
FROM pg_stat_user_tables
ORDER BY n_dead_tup DESC
LIMIT 20;
-- 查看 VACUUM 进度
SELECT * FROM pg_stat_progress_vacuum;
-- 排查阻碍 VACUUM 的长事务/会话
SELECT pid, state, xact_start, query
FROM pg_stat_activity
WHERE state <> 'idle'
ORDER BY xact_start;
长事务、复制槽滞留等都会阻止清理/冻结,导致表膨胀。
4) Prefect 侧的“止血”
- 缩短事件保留期:`PREFECT_SERVER_EVENTS_RETENTION_PERIOD="P15D"`
- 关闭/降低日志上报:`PREFECT_LOGGING_TO_API_ENABLED=false` 或 PREFECT_LOGGING_LEVEL=WARNING
结论
- autovacuum 会自动清理,但无法保证在你这种一次性大量删除后“及时”和“物理缩小文件”。因此仍建议:大删后手动 `VACUUM (ANALYZE)`;需要立即降磁盘用 `pg_repack`/`VACUUM FULL`;同时把 autovacuum 调得更积极,尤其是 events/log 表。
Prefect 的维护参考:Database maintenanceR G
11/11/2025, 3:04 AMMarvin
11/11/2025, 3:04 AMMarvin
11/11/2025, 3:09 AMR G
11/11/2025, 3:21 AMMarvin
11/11/2025, 3:21 AMMarvin
11/11/2025, 3:24 AMBEGIN;
-- 基准时间点(15天)
-- 可改:interval '15 days'
WITH cutoff AS (
SELECT NOW() - interval '15 days' AS ts
),
-- 需要删除的 flow_run
old_flow_runs AS (
SELECT fr.id
FROM flow_run fr, cutoff c
WHERE COALESCE(fr.end_time, fr.created) < c.ts
),
-- 需要删除的 task_run(隶属于 old_flow_runs)
old_task_runs AS (
SELECT tr.id
FROM task_run tr
JOIN old_flow_runs fr ON fr.id = tr.flow_run_id
),
-- 1) 解除可能的外键引用,避免删除冲突
-- 如果有 task_run 指向将被删除的子流(child_flow_run_id),先置空
nulled_child AS (
UPDATE task_run t
SET child_flow_run_id = NULL
WHERE t.child_flow_run_id IN (SELECT id FROM old_flow_runs)
RETURNING 1
),
-- 如果有 task_run 指向将被删除的父 task(parent_task_run_id),先置空
nulled_parent AS (
UPDATE task_run t
SET parent_task_run_id = NULL
WHERE t.parent_task_run_id IN (SELECT id FROM old_task_runs)
RETURNING 1
),
-- flow/task 当前状态外键(state_id)先置空,避免状态表删除时触发约束
nulled_task_state_ptr AS (
UPDATE task_run t
SET state_id = NULL
WHERE t.id IN (SELECT id FROM old_task_runs)
RETURNING 1
),
nulled_flow_state_ptr AS (
UPDATE flow_run f
SET state_id = NULL
WHERE f.id IN (SELECT id FROM old_flow_runs)
RETURNING 1
)
-- 2) 先删从表数据(状态、日志、工件、输入等)
-- task_run_state
DELETE FROM task_run_state s
USING old_task_runs r
WHERE s.task_run_id = r.id;
-- flow_run_state
DELETE FROM flow_run_state s
USING old_flow_runs r
WHERE s.flow_run_id = r.id;
-- 日志(可能很多,注意执行时间)
DELETE FROM "log" l
WHERE l.flow_run_id IN (SELECT id FROM old_flow_runs)
OR l.task_run_id IN (SELECT id FROM old_task_runs);
-- 工件
DELETE FROM artifact a
WHERE a.flow_run_id IN (SELECT id FROM old_flow_runs)
OR a.task_run_id IN (SELECT id FROM old_task_runs);
-- flow_run 输入
DELETE FROM flow_run_input i
USING old_flow_runs r
WHERE i.flow_run_id = r.id;
-- 如有其他引用 run 的表(例如自定义扩展),请在此位置追加 DELETE
-- 3) 删除主表
DELETE FROM task_run t
USING old_task_runs r
WHERE t.id = r.id;
DELETE FROM flow_run f
USING old_flow_runs r
WHERE f.id = r.id;
COMMIT;
二、批处理版(建议用于生产,控制单次删除量,降低锁表和膨胀)
- 创建一个过程,每次按 batch_size 分批清理
- 默认保留 15 天,可修改参数
- 可配合 pg_cron 定时执行
```
CREATE OR REPLACE PROCEDURE prefect_cleanup(retention_days integer DEFAULT 15, batch_size integer DEFAULT 5000)
LANGUAGE plpgsql
AS $$
DECLARE
cutoff_ts timestamptz := NOW() - make_interval(days => retention_days);
v_deleted bigint;
BEGIN
RAISE NOTICE 'Prefect cleanup start, cutoff=%', cutoff_ts;
-- 多轮循环,直到没有可删数据
LOOP
-- 解除外键:child_flow_run_id
WITH to_delete_flow_runs AS (
SELECT id
FROM flow_run
WHERE COALESCE(end_time, created) < cutoff_ts
LIMIT batch_size
)
UPDATE task_run t
SET child_flow_run_id = NULL
WHERE t.child_flow_run_id IN (SELECT id FROM to_delete_flow_runs);
GET DIAGNOSTICS v_deleted = ROW_COUNT;
-- 解除外键:parent_task_run_id
WITH to_delete_task_runs AS (
SELECT tr.id
FROM task_run tr
JOIN flow_run fr ON fr.id = tr.flow_run_id
WHERE COALESCE(fr.end_time, fr.created) < cutoff_ts
LIMIT batch_size
)
UPDATE task_run t
SET parent_task_run_id = NULL
WHERE t.parent_task_run_id IN (SELECT id FROM to_delete_task_runs);
GET DIAGNOSTICS v_deleted = v_deleted + ROW_COUNT;
-- 置空 state_id 指针(task)
WITH to_delete_task_runs AS (
SELECT tr.id
FROM task_run tr
JOIN flow_run fr ON fr.id = tr.flow_run_id
WHERE COALESCE(fr.end_time, fr.created) < cutoff_ts
LIMIT batch_size
)
UPDATE task_run t
SET state_id = NULL
WHERE t.id IN (SELECT id FROM to_delete_task_runs);
GET DIAGNOSTICS v_deleted = v_deleted + ROW_COUNT;
-- 置空 state_id 指针(flow)
WITH to_delete_flow_runs AS (
SELECT id
FROM flow_run
WHERE COALESCE(end_time, created) < cutoff_ts
LIMIT batch_size
)
UPDATE flow_run fMarvin
11/11/2025, 3:24 AM可选:用 pg_cron 定时
-- 安装扩展(如未启用)
-- CREATE EXTENSION IF NOT EXISTS pg_cron;
-- 每天凌晨 3:15 执行一次
SELECT cron.schedule(
'prefect_cleanup_daily',
'15 3 * * *',
$$CALL prefect_cleanup(15, 5000);$$
);
```
注意事项
- 强烈建议先做“干跑”验证:把 DELETE 改成 SELECT COUNT(*) FROM ... WHERE ... 看看涉及行数是否符合预期。
- 大表删除后请留意 Autovacuum;有需要可配合 pg_repack 做空间回收。
- 日志表行数可能极多,确保有索引:"log"(flow_run_id), "log"(task_run_id);Prefect 默认迁移会建索引。
- 事件表(event)通常不通过 FK 直接关联 run,且 Prefect 有事件保留策略;本脚本未清理 event,如需统一保留策略可单独处理。
参考
- Prefect 数据库维护指南
- ORM 模型(表结构参考)
如果你能提供具体的 Prefect/DB 版本或自定义表名,我可以帮你把脚本对齐到你的精确 schema。R G
11/11/2025, 11:31 AMWITH old_flows AS (
SELECT id FROM flow_run
WHERE end_time IS NOT NULL
AND end_time < now() - interval '15 days'
LIMIT 30000 -- 分批
)
DELETE FROM task_run_state
WHERE task_run_id IN (
SELECT id FROM task_run
WHERE flow_run_id IN (SELECT id FROM old_flows)
);
这个清理task_run_state有没有优化方法,执行很慢Marvin
11/11/2025, 11:31 AMMarvin
11/11/2025, 11:32 AMCREATE INDEX CONCURRENTLY IF NOT EXISTS idx_flow_run_end_time_notnull
ON flow_run (end_time)
WHERE end_time IS NOT NULL;
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_task_run_flow_run_id
ON task_run (flow_run_id);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_task_run_state_task_run_id
ON task_run_state (task_run_id);
推荐的删除写法(用 JOIN,按 task_run 分批)
WITH cutoff AS (
SELECT now() - interval '15 days' AS ts
),
to_delete_task_runs AS (
SELECT tr.id
FROM task_run tr
JOIN flow_run fr ON fr.id = tr.flow_run_id
JOIN cutoff c ON true
WHERE fr.end_time IS NOT NULL
AND fr.end_time < c.ts
ORDER BY tr.id
LIMIT 30000 -- 分批删除 3万条 task_run 对应的状态
)
DELETE FROM task_run_state s
USING to_delete_task_runs d
WHERE s.task_run_id = d.id;
为什么这样更快
- 避免了“task_run_id IN (SELECT … IN (SELECT …))”的嵌套子查询,方便优化器走索引
- 直接以 task_run 为粒度批量,能明确控制每批删除的状态行数
- 依赖的索引能被稳定命中:flow_run(end_time) → task_run(flow_run_id) → task_run_state(task_run_id)
如果你要进一步压低开销,可用临时表分两步(更清晰的执行计划、可重复使用)
-- 第一步:预计算本批要删的 task_run
CREATE TEMP TABLE tmp_old_task_runs ON COMMIT DROP AS
SELECT tr.id
FROM task_run tr
JOIN flow_run fr ON fr.id = tr.flow_run_id
WHERE fr.end_time IS NOT NULL
AND fr.end_time < now() - interval '15 days'
ORDER BY tr.id
LIMIT 30000;
-- 第二步:基于临时表删除状态
DELETE FROM task_run_state s
USING tmp_old_task_runs t
WHERE s.task_run_id = t.id;
执行与维护建议
- 先用 EXPLAIN (ANALYZE, BUFFERS) 对比你的原语句和上面的重写,确认索引命中和 I/O 降低
- 每批删除后跑
ANALYZE task_run_state;
- 大量清理完成后,交给 autovacuum,必要时使用 pg_repack 回收表膨胀
- 多次循环批量执行,直到本批删除行数接近 0
补充:如果你后续也要删除 task_run/flow_run 本体,务必先置空可能的外键指针(如 task_run.state_id、flow_run.state_id、task_run.child_flow_run_id、task_run.parent_task_run_id),再删从表(状态、日志、工件等),最后删主表;我可以给你一整套分批过程化脚本。
参考
- Prefect 数据库维护指南R G
11/11/2025, 11:41 AMMarvin
11/11/2025, 11:41 AMMarvin
11/11/2025, 11:42 AMBEGIN;
-- 1) 冻结本批要删除的 flow_run(15 天前完成)
CREATE TEMP TABLE _del_flow_runs ON COMMIT DROP AS
SELECT id
FROM flow_run
WHERE end_time IS NOT NULL
AND end_time < now() - interval '15 days'
ORDER BY id
LIMIT 30000; -- 分批量
-- 没有可删则退出
DO $$
BEGIN
IF (SELECT COUNT(*) FROM _del_flow_runs) = 0 THEN
RAISE NOTICE 'no flow_runs to delete';
END IF;
END$$;
-- 2) 冻结本批相关的 task_run
CREATE TEMP TABLE _del_task_runs ON COMMIT DROP AS
SELECT tr.id
FROM task_run tr
JOIN _del_flow_runs d ON d.id = tr.flow_run_id;
-- 3) 先解除可能的外键指针,避免约束冲突(若你的 schema 无这些列,可略过)
UPDATE task_run t
SET child_flow_run_id = NULL
WHERE t.child_flow_run_id IN (SELECT id FROM _del_flow_runs);
UPDATE task_run t
SET parent_task_run_id = NULL
WHERE t.parent_task_run_id IN (SELECT id FROM _del_task_runs);
UPDATE task_run t
SET state_id = NULL
WHERE t.id IN (SELECT id FROM _del_task_runs);
UPDATE flow_run f
SET state_id = NULL
WHERE f.id IN (SELECT id FROM _del_flow_runs);
-- 4) 删从表(状态、日志、工件、输入等)
DELETE FROM task_run_state s
USING _del_task_runs t
WHERE s.task_run_id = t.id;
DELETE FROM flow_run_state s
USING _del_flow_runs f
WHERE s.flow_run_id = f.id;
DELETE FROM "log" l
WHERE l.flow_run_id IN (SELECT id FROM _del_flow_runs)
OR l.task_run_id IN (SELECT id FROM _del_task_runs);
DELETE FROM artifact a
WHERE a.flow_run_id IN (SELECT id FROM _del_flow_runs)
OR a.task_run_id IN (SELECT id FROM _del_task_runs);
DELETE FROM flow_run_input i
USING _del_flow_runs f
WHERE i.flow_run_id = f.id;
-- 5) 删主表
DELETE FROM task_run t
USING _del_task_runs dt
WHERE t.id = dt.id;
DELETE FROM flow_run f
USING _del_flow_runs df
WHERE f.id = df.id;
COMMIT;
加固建议
- 始终在一个事务里完成“冻结集合 → 删除”,本批内不会因为“时间差”而遗漏。
- 按上述顺序删除(子→父),避免非级联外键导致的报错。
- 如需进一步避免和其他清理任务抢相同记录,可在冻结集合时考虑行级锁(例如用一个标记表记录 id,或用 FOR UPDATE SKIP LOCKED 选取需要的行;不过 FOR UPDATE 并不总能直接用于 CTAS)。
- 用循环或定时器重复执行该事务,直到每批返回 0 行。
- 大删后 ANALYZE 相关表;空间回收可交给 autovacuum,必要时用 pg_repack。
如果你想直接依赖数据库级的级联删除,也可以先确认外键的 delete_rule:
SELECT
tc.table_name, ccu.table_name AS referenced_table, rc.delete_rule
FROM information_schema.table_constraints AS tc
JOIN information_schema.referential_constraints rc
ON rc.constraint_name = tc.constraint_name AND rc.constraint_schema = tc.table_schema
JOIN information_schema.constraint_column_usage AS ccu
ON ccu.constraint_name = tc.constraint_name AND ccu.constraint_schema = tc.table_schema
WHERE tc.constraint_type = 'FOREIGN KEY'
AND tc.table_schema = 'public'
AND (tc.table_name, ccu.table_name) IN (
('task_run','flow_run'),
('task_run_state','task_run'),
('flow_run_state','flow_run')
);
- 如果看到 delete_rule = CASCADE,则可以只删父表,数据库会级联清理子表。
- 如果是 NO ACTION/RESTRICT,则必须按上面的顺序手动删子再删父。R G
11/11/2025, 11:48 AMMarvin
11/11/2025, 11:48 AMMarvin
11/11/2025, 11:50 AM-- 查关键外键的 delete_rule(CASCADE / SET NULL / NO ACTION)
SELECT
tc.table_name AS fk_table,
kcu.column_name AS fk_column,
ccu.table_name AS pk_table,
rc.delete_rule
FROM information_schema.table_constraints AS tc
JOIN information_schema.key_column_usage AS kcu
ON kcu.constraint_name = tc.constraint_name AND kcu.constraint_schema = tc.constraint_schema
JOIN information_schema.referential_constraints AS rc
ON rc.constraint_name = tc.constraint_name AND rc.constraint_schema = tc.constraint_schema
JOIN information_schema.constraint_column_usage AS ccu
ON ccu.constraint_name = tc.constraint_name AND ccu.constraint_schema = tc.constraint_schema
WHERE tc.constraint_type = 'FOREIGN KEY'
AND tc.table_schema = 'public'
AND (tc.table_name, kcu.column_name) IN (
('task_run','flow_run_id'),
('task_run','parent_task_run_id'),
('task_run','child_flow_run_id'),
('task_run','state_id'),
('flow_run','state_id'),
('task_run_state','task_run_id'),
('flow_run_state','flow_run_id')
)
ORDER BY fk_table, fk_column;
如何据此决定是否可以省略“置空外键指针”
- 如果看到:
- task_run_state.task_run_id → ON DELETE CASCADE
- flow_run_state.flow_run_id → ON DELETE CASCADE
- task_run.state_id 和 flow_run.state_id → ON DELETE SET NULL
- task_run.child_flow_run_id、parent_task_run_id → ON DELETE SET NULL 或你能保证同批一起删掉被引用的行
那么你可以不手动置空指针,直接按“先删子表(或直接删父表触发级联)→ 再删父表”的流程。
- 只要其中任何一条不是 CASCADE/SET NULL,就不要省略置空;否则会被外键阻止删除。
更安全/通用的做法(即使外键不是 CASCADE/SET NULL 也能用)
- 保留“置空外键指针”这一步,但确保只对“本批要删的集合”做更新(用临时表冻结集合),这样额外开销有限。
- 关键列加索引,加速置空步骤中的 WHERE 过滤:
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_task_run_child_flow_run_id ON task_run (child_flow_run_id);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_task_run_parent_task_run_id ON task_run (parent_task_run_id);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_task_run_state_id ON task_run (state_id);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_flow_run_state_id ON flow_run (state_id);
如果你想完全不置空指针,还有两条路可选(需你确认/变更架构):
- 依赖级联:保证上述外键是 CASCADE 或 SET NULL,然后按“删父触发子级联”的顺序删除。
- 调整删除顺序与过滤逻辑:例如删除 state 表时跳过“当前状态”那一行(被 state_id 引用),在删完 run 后再补删剩余 state。但实现起来更复杂、也容易遗漏。
建议
- Prefect 自托管场景下,默认不保证这些外键是 CASCADE/SET NULL。因此,除非你查到 delete_rule 已满足,否则不要省略置空步骤。
- 如果你贴一下查询结果(delete_rule 列),我可以给你一个最简版的“无需置空/或仅置空某几列”的脚本。
参考
- Prefect 数据库维护指南R G
11/11/2025, 11:52 AMMarvin
11/11/2025, 11:53 AMMarvin
11/11/2025, 11:53 AM-- 准备索引(只需建一次;没有这个索引会很慢)
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_flow_run_end_time_notnull
ON flow_run (end_time) WHERE end_time IS NOT NULL;
-- 按批次删除 3 万个超过 15 天的 flow_run,其它表走级联
WITH del AS (
SELECT id
FROM flow_run
WHERE end_time IS NOT NULL
AND end_time < now() - interval '15 days'
ORDER BY id
LIMIT 30000
)
DELETE FROM flow_run f
USING del d
WHERE f.id = d.id;
-- 可循环执行,直到受影响行数为 0
注意
- 你发的外键列表里没有涉及日志/工件等表。如果这些表没有对 flow_run_id/task_run_id 建立带 CASCADE 的外键,就不会自动级联,需要单独清理或给它们加外键/触发器。
- 如果需要,我可以给你补充 logs/artifact/flow_run_input 的清理批处理。
- 大量删除后跑 ANALYZE,空间回收交给 autovacuum;需要强制回收时用 pg_repack,而不是 VACUUM FULL。R G
11/11/2025, 11:59 AMMarvin
11/11/2025, 11:59 AMMarvin
11/11/2025, 12:01 PM-- 按 end_time 过滤 flow_run(部分索引,避免扫描 NULL)
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_flow_run_end_time_notnull
ON flow_run (end_time) WHERE end_time IS NOT NULL;
-- 连接 task_run → flow_run(有外键也建议有索引)
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_task_run_flow_run_id
ON task_run (flow_run_id);
-- 日志按 run 关联删除(强烈建议)
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_flow_run_id
ON "log" (flow_run_id);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_task_run_id
ON "log" (task_run_id);
-- 可选:工件、flow_run_input(若存在这些表,建议加索引)
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_artifact_flow_run_id
ON artifact (flow_run_id);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_artifact_task_run_id
ON artifact (task_run_id);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_flow_run_input_flow_run_id
ON flow_run_input (flow_run_id);
二、单次“分批删除”的事务脚本(默认每批 3 万个 flow_run)
- 保留最近 15 天:end_time < now() - interval '15 days'
- 日志按 run 关联清理
- 依赖你贴出的外键,task_run + 两类 state 会自动级联删除
BEGIN;
-- 1) 冻结本批要删除的 flow_run
CREATE TEMP TABLE _del_flow_runs ON COMMIT DROP AS
SELECT id
FROM flow_run
WHERE end_time IS NOT NULL
AND end_time < now() - interval '15 days'
ORDER BY id
LIMIT 30000; -- 调整批量大小以平衡时长与锁
-- 2) 冻结本批关联的 task_run(用于日志/工件等非级联表清理)
CREATE TEMP TABLE _del_task_runs ON COMMIT DROP AS
SELECT tr.id
FROM task_run tr
JOIN _del_flow_runs d ON d.id = tr.flow_run_id;
-- 3) 先清理日志(非级联)
DELETE FROM "log" l
WHERE l.flow_run_id IN (SELECT id FROM _del_flow_runs)
OR l.task_run_id IN (SELECT id FROM _del_task_runs);
-- 4) 可选:清理工件(如存在该表)
DELETE FROM artifact a
WHERE a.flow_run_id IN (SELECT id FROM _del_flow_runs)
OR a.task_run_id IN (SELECT id FROM _del_task_runs);
-- 5) 可选:清理 flow_run_input(如存在该表)
DELETE FROM flow_run_input i
USING _del_flow_runs d
WHERE i.flow_run_id = d.id;
-- 6) 删除 flow_run(将自动级联删除 task_run / flow_run_state / task_run_state)
DELETE FROM flow_run f
USING _del_flow_runs d
WHERE f.id = d.id;
COMMIT;
使用方式
- 将上述“单次分批删除”脚本周期性执行,直到本批删除行数为 0 即完成本轮清理。
- 如果你希望自动化执行,可以使用 pg_cron 每晚跑多次。例如每晚跑 10 次,每次 3 万条,具体次数/批量按你数据量和时间窗调节。
pg_cron 示例(可选)
-- 如未启用
-- CREATE EXTENSION IF NOT EXISTS pg_cron;
-- 每天 03:00 开始,每 10 分钟跑一次,共 6 次(示意)
SELECT cron.schedule('prefect_cleanup_03_00', '0 3 * * *', $$BEGIN;
CREATE TEMP TABLE _del_flow_runs ON COMMIT DROP AS
SELECT id FROM flow_run WHERE end_time IS NOT NULL AND end_time < now() - interval '15 days'
ORDER BY id LIMIT 30000;
CREATE TEMP TABLE _del_task_runs ON COMMIT DROP AS
SELECT tr.id FROM task_run tr JOIN _del_flow_runs d ON d.id = tr.flow_run_id;
DELETE FROM "log" l WHERE l.flow_run_id IN (SELECT id FROM _del_flow_runs) OR l.task_run_id IN (SELECT id FROM _del_task_runs);
DELETE FROM artifact a WHERE a.flow_run_id IN (SELECT id FROM _del_flow_runs) OR a.task_run_id IN (SELECT id FROM _del_task_runs);
DELETE FROM flow_run_input i USING _del_flow_runs d WHERE i.flow_run_id = d.id;
DELETE FROM flow_run f USING _del_flow_runs d WHERE f.id = d.id;
COMMIT;$$);
SELECT cron.schedule('prefect_cleanup_03_10', '10 3 * * *', $$...同上脚本...$$);
-- 依次 03:20、03:30、03:40、03:50 再调度几次
可选:日志按“时间”统一保留 15 天
- 如果你希望日志本身也不超过 15 天(即使其 run 仍未被删除),需要在 "log" 表上按“时间列”做时间删除并配相应索引。
- Prefect 的日志表时间列名请以你的实际 schema 为准(有的环境是 created,有的可能是 timestamp)。建索引后再按时间删除:
-- 假设日志表时间列名为 created(请先确认你的列名)
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_created ON "log"(created);
DELETE FROM "log"
WHERE created < now() - interval '15 days'
LIMIT 500000; -- 可选:分批删除,具体可用 CTE + ORDER BY + LIMIT
维护建议
- 大量删除后,执行
```
ANALYZE flow_run;
ANALYZE task_run;Marvin
11/11/2025, 12:01 PMR G
11/12/2025, 9:20 AMMarvin
11/12/2025, 9:20 AMMarvin
11/12/2025, 9:21 AM-- 两个都试一下,哪个不是 NULL 就说明实际注册的表名
SELECT to_regclass('public.log');
SELECT to_regclass('public."log"');
-- psql 里也可以
-- \d log
-- \d "log"
结论与建议
- Prefect 3.x 默认是小写 log 表,你可以直接用:
DELETE FROM log l
WHERE l.flow_run_id IN (SELECT id FROM _del_flow_runs)
OR l.task_run_id IN (SELECT id FROM _del_task_runs);
- 如果你的环境确实是大小写或特殊命名,继续保留双引号即可。R G
11/12/2025, 10:56 AMMarvin
11/12/2025, 10:56 AMMarvin
11/12/2025, 10:57 AMCREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_flow_run_id_notnull
ON log (flow_run_id) WHERE flow_run_id IS NOT NULL;
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_task_run_id_notnull
ON log (task_run_id) WHERE task_run_id IS NOT NULL;
2) 把一条“带 OR 的删除”改成两条 JOIN 删除(避免 OR)
替换你脚本里原先的
DELETE FROM log l
WHERE l.flow_run_id IN (SELECT id FROM _del_flow_runs)
OR l.task_run_id IN (SELECT id FROM _del_task_runs);
为下面两条(更容易命中索引):
-- 先按 flow_run_id 删
DELETE FROM log l
USING _del_flow_runs d
WHERE l.flow_run_id = d.id;
-- 再按 task_run_id 删
DELETE FROM log l
USING _del_task_runs t
WHERE l.task_run_id = t.id;
说明:如果某些日志同时拥有 flow_run_id 和 task_run_id,第一条就会删掉它们,第二条不会重复删除(已不再存在)。
3) 如果仍慢:用 ctid 微批删除,避免一次性大范围扫描
把删除 log 的步骤替换为以下循环(每批 10 万行,你可按情况调小到 20k/50k):
DO $$
DECLARE
batch_limit integer := 100000;
rows_deleted integer;
BEGIN
-- 先针对 flow_run_id 的日志
LOOP
WITH to_del AS (
SELECT l.ctid
FROM log l
JOIN _del_flow_runs d ON l.flow_run_id = d.id
LIMIT batch_limit
)
DELETE FROM log WHERE ctid IN (SELECT ctid FROM to_del);
GET DIAGNOSTICS rows_deleted = ROW_COUNT;
EXIT WHEN rows_deleted = 0;
END LOOP;
-- 再针对 task_run_id 的日志
LOOP
WITH to_del AS (
SELECT l.ctid
FROM log l
JOIN _del_task_runs t ON l.task_run_id = t.id
LIMIT batch_limit
)
DELETE FROM log WHERE ctid IN (SELECT ctid FROM to_del);
GET DIAGNOSTICS rows_deleted = ROW_COUNT;
EXIT WHEN rows_deleted = 0;
END LOOP;
END$$;
如果日志量极大,建议把第 3 步放到事务外循环执行(每批提交一次),进一步降低锁持有时长。
4) 如果你的策略是“日志也只保留 15 天”,可以改为纯时间清理(通常更快)
前提:给日志时间列建索引(例如 created;请确认你库里的列名)
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_created ON log (created);
-- 微批按时间删(Postgres 不允许 DELETE … ORDER BY)
DO $$
DECLARE
rows_deleted integer;
BEGIN
LOOP
WITH to_del AS (
SELECT ctid
FROM log
WHERE created < now() - interval '15 days'
LIMIT 100000
)
DELETE FROM log WHERE ctid IN (SELECT ctid FROM to_del);
GET DIAGNOSTICS rows_deleted = ROW_COUNT;
EXIT WHEN rows_deleted = 0;
END LOOP;
END$$;
这样即便不关联 flow_run/_del_task_runs,也能把超过保留期的日志快速滚动清掉。
5) 维护与观察
- 大批量删除后,执行 ANALYZE,留给 autovacuum 做清理;需要强制回收用 pg_repack
ANALYZE log;
- 避免“长事务占着旧快照”拖慢 autovacuum:大批删除建议切成多次短事务
- 如果执行计划依然走顺序扫描,检查统计信息是否过期、索引是否被正确选择;必要时先 EXPLAIN (ANALYZE, BUFFERS) 对比“OR 写法”和“两条 JOIN 写法”的差异
把以上第 2 步替换进你原先的事务脚本通常就能显著提速;如果数据量特别大,再叠加第 3 步的 ctid 微批。需要我帮你把整段事务脚本改成“ctid 微批 + 多次提交”的版本,也可以直接贴出来。R G
11/13/2025, 7:37 AMMarvin
11/13/2025, 7:37 AMMarvin
11/13/2025, 7:42 AMmy_task() 时,异常会像普通 Python 函数一样同步抛出;你可以在 flow 里用一个大的 try/except 包住整段逻辑。
from prefect import flow, task, get_run_logger
@task
def may_fail(x):
if x == 2:
raise RuntimeError(f"blade_req_final_failed=x={x}")
return x * 10
@flow
def my_flow():
logger = get_run_logger()
try:
a = may_fail(1) # 正常
b = may_fail(2) # 这里抛出异常,直接被下面的 except 捕获
c = may_fail(3) # 不会执行到
except Exception as exc:
logger.error(f"Flow 捕获到 task 异常: {exc}", exc_info=True)
# 需要让 Flow 标记为失败就 re-raise;只记录不失败就不 raise
raise
if __name__ == "__main__":
my_flow()
二、任务并发(用 submit 返回 PrefectFuture)
- 用 task.submit(...) 时,返回的是 `PrefectFuture`;异常不会立刻抛出,而是在你获取结果或检查 state 时体现。你可以:
- 通过 future.result() 在 flow 里用 try/except 捕获
- 或者先 wait/as_completed 等所有任务结束,再检查失败的任务 state 并打印异常
1) 最简单:逐个 result,并在 flow 顶层 try/except 中捕获任意一个失败
from prefect import flow, task, get_run_logger
from prefect.futures import as_completed
@task
def may_fail(x):
if x == 2:
raise RuntimeError(f"blade_req_final_failed=x={x}")
return x * 10
@flow
def my_flow():
logger = get_run_logger()
futs = [may_fail.submit(i) for i in range(5)]
try:
# 任意一个 future 在 .result() 时失败,都会被这里的 except 捕获
for fut in as_completed(futs):
_ = fut.result() # 默认 raise_on_failure=True
except Exception as exc:
logger.error(f"Flow 捕获到 task 异常: {exc}", exc_info=True)
# 想让 Flow 失败就继续 raise
raise
2) 想“全部执行完后一次性打印所有失败”的异常细节(不只第一个):
from prefect import flow, task, get_run_logger
from prefect.futures import wait
from prefect.states import get_state_exception
@task
def may_fail(x):
if x % 3 == 0:
raise RuntimeError(f"blade_req_final_failed=x={x}")
return x * 10
@flow
def my_flow():
logger = get_run_logger()
futs = [may_fail.submit(i) for i in range(1, 10)]
# 等待所有 future 完成
done, _ = wait(futs)
failed = []
for fut in done:
st = fut.state
if st.is_failed():
exc = get_state_exception(st) # 拿到原始异常
logger.error(f"Task 失败: task_run_id={fut.task_run_id}, exc={exc}", exc_info=True)
failed.append(exc)
if failed:
# 如需让 Flow 标记为失败,可以合并/抛出一个代表性异常
raise RuntimeError(f"{len(failed)} task(s) failed; see logs above.")
要点
- 直接调用 task:异常同步抛出 -> 用普通 try/except。
- 并发提交 task:`submit()` 返回 Future -> 在 flow 里用 `future.result()`、`wait`/`as_completed`、或检查 future.state + get_state_exception(state) 来捕获/打印。
- 在 flow 里用 get_run_logger() 打日志,方便在 Prefect UI 中查看。
- 如果你想在 task 自身失败时自动回调打印日志,也可以给 task 配置 on_failure 钩子,但你当前场景只在 flow 捕获即可。
参考
- Prefect 3 - States 概念
- PrefectFutures API
- 并发运行任务
- 在 Flow/Task 中记录日志
如果你现在是直接调用还是用 .submit() 并发?我可以根据你的具体调用方式帮你把示例改到最贴合你的代码。R G
11/13/2025, 7:50 AMMarvin
11/13/2025, 7:51 AMMarvin
11/13/2025, 7:51 AMsubmit() 并发提交;在 flow 顶层用 try/except 做总兜底
- 用 wait() 等所有已提交的 futures 完成,并通过 future.state 判断失败是谁,打印详细异常
- 只在需要时调用 `.result()`(比如你的条件判断和最终返回)
from prefect import flow, get_run_logger
from prefect.context import get_run_context
from prefect.futures import wait
from prefect.states import get_state_exception
import traceback
@flow(flow_run_name="{logid}")
def my_flow(input_data, logid):
logger = get_run_logger()
mario_log.notice(f"act=my_flow_run, logid={logid}")
flow_run_id = get_run_context().flow_run.id
try:
# 并发提交
a = task_a.submit(input_data)
b = task_b.with_options(retries=3).submit(a)
c = task_c.submit(a)
tasks = {"b": b, "c": c}
# 用 a 的结果做条件判断;如果 a 失败,这里会抛出并被外层 except 捕获
a_value = a.result()
if a_value["value"] > 30:
e = task_e.submit(a)
tasks["e"] = e
d = task_d.submit(tasks)
# 等待所有已提交任务完成,并逐个检查失败,打印具体是哪个 task 挂了
name_to_future = {"a": a, **tasks, "d": d}
wait(list(name_to_future.values()))
any_failed = False
for name, fut in name_to_future.items():
st = fut.state
if st.is_failed():
exc = get_state_exception(st)
mario_log.warning(f"act=task_failed, logid={logid}, task={name}, exc={exc}")
logger.error(f"task={name} failed: {exc}", exc_info=True)
any_failed = True
if any_failed:
# 统一走到 except 做失败返回
raise RuntimeError("One or more tasks failed, see logs above.")
# 所有任务成功,返回最终结果
return {"flow_run_id": flow_run_id, "errno": 0, "errmsg": "", "result": d.result()}
except Exception as e:
mario_log.warning(f"act=run_flow_failed, logid={logid}, traceback={traceback.format_exc()}")
logger.error(f"Flow failed: {e}", exc_info=True)
return {"flow_run_id": flow_run_id, "errno": 1, "errmsg": str(e)}
说明
- wait(...) 不会抛异常,只是等待完成;我们用 future.state.is_failed() + get_state_exception(state) 来打印每个失败的 task 的异常细节。
- 你的条件判断保持不变:`a.result()` 如果失败,会在这里抛出并被外层 except 捕获。
- 需要获取 flow_run_id 时,用 `get_run_context().flow_run.id`。
文档参考
- PrefectFutures API
- States 概念
如果你希望“不在条件判断里提前抛出 a 的异常”,也可以先检查 a.state 再决定是否访问结果,我可以再给你改一版。R G
11/13/2025, 7:55 AMMarvin
11/13/2025, 7:55 AMMarvin
11/13/2025, 7:56 AMd 调一次 `.result()`。如果任一上游 task(a/b/c/e/d 里任何一个)失败,Prefect 会在 d.result() 时抛出原始异常或 UpstreamTaskError,被你外层的 try/except 统一捕获并打印。
关键点:把对 a 的条件判断从 flow 移进 task_e 自身,这样在 flow 中完全不需要提前 `a.result()`。示例改造如下:
from prefect import flow
from prefect.context import get_run_context
import traceback
@flow(flow_run_name="{logid}")
def my_flow(input_data, logid):
""" my_flow """
mario_log.notice(f"act=my_flow_run, logid={logid}")
flow_run_id = get_run_context().flow_run.id
try:
a = task_a.submit(input_data)
b = task_b.with_options(retries=3).submit(a)
c = task_c.submit(a)
# 关键:总是提交 e,但把“是否需要执行”的判断放进 task_e 里
# 如果不满足条件,task_e 自己直接 no-op/返回 None
e = task_e.submit(a)
tasks = {'b': b, 'c': c, 'e': e}
d = task_d.submit(tasks)
# 只在最后取一次结果;任何上游失败都会在这里抛出异常
result = d.result()
return {"flow_run_id": flow_run_id, "errno": 0, "errmsg": "", "result": result}
except Exception as e:
mario_log.warning(f"act=run_flow_failed, logid={logid}, traceback={traceback.format_exc()}")
return {"flow_run_id": flow_run_id, "errno": 1, "errmsg": str(e)}
对应的 task_e 请改成内部自判是否需要执行,避免在 flow 里先 `.result()`:
from prefect import task
@task
def task_e(a_value):
# a_value 会由 Prefect 自动解析为 task_a 的结果
if a_value["value"] <= 30:
return None # 不需要执行时直接返回
# 否则执行你的原逻辑
# ... 原有 task_e 的实现 ...
为什么这样可以不逐个检查
- d 依赖 `tasks`(包含 b/c/e),这些又依赖 `a`。如果链路上任何一个失败,`d.result()` 会抛出异常,外层 try/except 统一捕获并打印。
- Flow 中没有任何其他 .result() 调用,自然就不存在“逐个检查”的负担。
可选等价写法(同样只需一次触发异常)
- 用一次性解析触发异常:
from prefect.futures import resolve_futures_to_results
resolve_futures_to_results([d]) # 任一上游失败会在这里抛出
result = d.result()
参考
- States 与上游失败传播
- PrefectFutures APIR G
11/13/2025, 8:00 AMMarvin
11/13/2025, 8:00 AMMarvin
11/13/2025, 8:05 AMd.result() 做一次 try/except
- 一旦捕获异常,就用 resolve_futures_to_states(...) 一次性把所有相关 futures 解析成 states
- 从失败的 state 中用 get_state_exception(state) 提取真实异常并打印
- 不需要在主逻辑里逐个 if 检查哪个 task 失败
改造代码如下(延续我们之前把 e 的条件判断移进 task_e 的思路,避免在 flow 中提前 `a.result()`):
from prefect import flow, get_run_logger
from prefect.context import get_run_context
from prefect.futures import resolve_futures_to_states
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
import traceback
# 可选:task_e 内部做条件判断,不满足则直接 no-op
# @task
# def task_e(a_value):
# if a_value["value"] <= 30:
# return None
# ... # 你的逻辑
@flow(flow_run_name="{logid}")
def my_flow(input_data, logid):
""" my_flow """
logger = get_run_logger()
mario_log.notice(f"act=my_flow_run, logid={logid}")
flow_run_id = get_run_context().flow_run.id
try:
a = task_a.submit(input_data)
b = task_b.with_options(retries=3).submit(a)
c = task_c.submit(a)
e = task_e.submit(a) # 是否需要执行在 task_e 内部判断
tasks = {"b": b, "c": c, "e": e}
d = task_d.submit(tasks)
# 只在这里取一次结果;任何上游失败都会在这里抛出
result = d.result()
return {"flow_run_id": flow_run_id, "errno": 0, "errmsg": "", "result": result}
except Exception as err:
# 统一失败日志
mario_log.warning(f"act=run_flow_failed, logid={logid}, traceback={traceback.format_exc()}")
logger.error(f"Flow failed: {err}", exc_info=True)
# 一次性解析所有相关 futures -> states,然后把失败根因全部打印出来
try:
name_to_future = {"a": a, **tasks, "d": d}
states = resolve_futures_to_states(name_to_future)
for name, st in states.items():
if st.is_failed():
exc = get_state_exception(st)
mario_log.warning(
f"act=task_failed, logid={logid}, task={name}, exc={exc}, tb={exception_traceback(exc)}"
)
logger.error(f"task={name} failed: {exc}\n{exception_traceback(exc)}")
except Exception:
# 即使解析/打印失败也不影响主返回
pass
return {"flow_run_id": flow_run_id, "errno": 1, "errmsg": str(err)}
说明
- UpstreamTaskError 不携带根因异常;通过 resolve_futures_to_states(...) + get_state_exception(state) 才能拿到真正的 raise 信息和 traceback。
- 这段代码没有在主流程里“逐个检查 task 是否失败”,只是在最终失败时统一解析、统一打印。
- 如果你完全不想在 flow 里做任何遍历,推荐给各 task 配置 on_failure 钩子,在 task 失败的瞬间自动记录根因日志;flow 里只需捕获一次 UpstreamTaskError 即可。参考:State change hooks。
更多参考
- PrefectFutures API
- States 概念
- 在 Flow/Task 中记录日志R G
11/13/2025, 8:06 AMMarvin
11/13/2025, 8:06 AMMarvin
11/13/2025, 8:08 AMget_state_exception(state) 取到 task/flow 的真实异常;`state.message` 则是 Prefect 格式化后的错误信息。
- 日志:建议用 `get_run_logger()`,也可用你已有的 `mario_log`。
一、给 task 加 on_failure
- 用装饰器或 with_options 都可以。你可以一个钩子复用在多个 task 上。
from prefect import task, flow, get_run_logger
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
def log_task_failure(task, task_run, state):
# 这里不会重试,说明已是最终失败
exc = get_state_exception(state)
logger = get_run_logger()
# 你也可以用 mario_log
logger.error(
f"[TASK FAILED] task={task.name} run_id={task_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
@task(on_failure=[log_task_failure])
def task_a(input_data):
...
# 如果你不想修改原装饰器,也可以在使用处挂钩子
# task_b = task_b.with_options(on_failure=[log_task_failure])
二、给 flow 加 on_failure(可选)
- 如果你希望 flow 失败时也自动打印根因,同样可以挂 flow 级别钩子。
from prefect.states import get_state_exception
def log_flow_failure(flow, flow_run, state):
exc = get_state_exception(state)
logger = get_run_logger()
logger.error(
f"[FLOW FAILED] flow={flow.name} flow_run_id={flow_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
@flow(flow_run_name="{logid}", on_failure=[log_flow_failure])
def my_flow(input_data, logid):
...
三、结合你的场景的完整示例(保持“只在最终取一次结果”)
- 把 task_e 的条件判断挪到任务内部,避免在 flow 里先 `.result()`。
- 当任一 task 失败时,该 task 的 on_failure 会自动打印出原始异常和 traceback;flow 里只在最后一次 d.result() 触发失败,不需要你手动逐个检查。
from prefect import flow, get_run_logger
from prefect.context import get_run_context
import traceback
def log_task_failure(task, task_run, state):
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
exc = get_state_exception(state)
logger = get_run_logger()
logger.error(
f"[TASK FAILED] task={task.name} run_id={task_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
# 如需同步打到你的日志系统:
# mario_log.warning(f"act=task_failed, logid={task_run.id}, task={task.name}, exc={exc}")
@task(on_failure=[log_task_failure])
def task_a(input_data):
...
@task(on_failure=[log_task_failure])
def task_b(a_value):
...
@task(on_failure=[log_task_failure])
def task_c(a_value):
...
@task(on_failure=[log_task_failure])
def task_e(a_value):
if a_value["value"] <= 30:
return None
# 否则执行你的逻辑
...
@task(on_failure=[log_task_failure])
def task_d(tasks_dict):
...
@flow(flow_run_name="{logid}")
def my_flow(input_data, logid):
logger = get_run_logger()
mario_log.notice(f"act=my_flow_run, logid={logid}")
flow_run_id = get_run_context().flow_run.id
try:
a = task_a.submit(input_data)
b = task_b.with_options(retries=3).submit(a)
c = task_c.submit(a)
e = task_e.submit(a)
tasks = {'b': b, 'c': c, 'e': e}
d = task_d.submit(tasks)
result = d.result() # 任一上游失败会在此抛出,上面的 on_failure 已经记录过细节
return {"flow_run_id": flow_run_id, "errno": 0, "errmsg": "", "result": result}
except Exception as e:
# Flow 级兜底日志(简单化即可,详细根因已由 task 的 on_failure 打过)
logger.error(f"Flow failed: {e}", exc_info=True)
mario_log.warning(f"act=run_flow_failed, logid={logid}, traceback={traceback.format_exc()}")
return {"flow_run_id": flow_run_id, "errno": 1, "errmsg": str(e)}
小结
- on_failure 会在最终失败时自动执行,你可以集中打印“原始异常”和完整 traceback。
- 优点:不需要在 flow 里逐个检查谁失败;每个失败 task 自己把异常说清楚。
- 如还需附加上下文(参数、外部标识等),可以在 hook 里读取 task_run.parameters 或在日志里带上 `logid`(可通过参数传递或读取 get_run_context() 补齐)。
文档参考
- State change hooksR G
11/13/2025, 8:17 AMMarvin
11/13/2025, 8:17 AMMarvin
11/13/2025, 8:18 AMprefect.task 预设默认参数,然后在项目里都用你重绑定后的 task 装饰器。
from functools import partial
from prefect import task as _task, get_run_logger
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
def log_task_failure(task, task_run, state):
exc = get_state_exception(state)
logger = get_run_logger()
logger.error(
f"[TASK FAILED] task={task.name} run_id={task_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
# 重绑定:给 task 预设默认的 on_failure 钩子
task = partial(_task, on_failure=[log_task_failure])
# 之后在本模块里,直接用 @task,就会自动带上 on_failure
@task
def task_a(x):
...
# 需要额外 options 也可以叠加传递
@task(retries=3, retry_delay_seconds=5)
def task_b(x):
...
建议把以上写在你项目的一个公共模块(比如 my_prefect.py),之后各处都 from my_prefect import task 使用。
方案 B:自定义一个“带默认钩子”的装饰器(可合并额外 on_failure)
- 更灵活:允许每个 task 额外再加自己的 on_failure,自动与默认钩子合并。
from prefect import task as _task, get_run_logger
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
def log_task_failure(task, task_run, state):
exc = get_state_exception(state)
logger = get_run_logger()
logger.error(
f"[TASK FAILED] task={task.name} run_id={task_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
def monitored_task(__fn=None, on_failure=None, **opts):
# 合并默认钩子 + 额外钩子
hooks = [log_task_failure] + (on_failure or [])
def _decorate(fn):
return _task(on_failure=hooks, **opts)(fn)
# 既支持 @monitored_task,也支持 @monitored_task(retries=3) 用法
return _decorate if __fn is None else _decorate(__fn)
# 用法示例
@monitored_task
def task_a(x):
...
@monitored_task(retries=3)
def task_b(x):
...
def extra_hook(task, task_run, state): ...
@monitored_task(on_failure=[extra_hook])
def task_c(x):
...
补充
- 对“已经定义好的 task”(比如三方包或别处定义)也可一次性加默认钩子:
- 在集中注册处做一次 .with_options(on_failure=[log_task_failure]) 包装,然后项目里都用包装后的引用即可。
- 例如:
from other.module import raw_task_b as _raw_task_b
task_b = _raw_task_b.with_options(on_failure=[log_task_failure])
- flow 同理可以默认化(如果需要):
from functools import partial
from prefect import flow as _flow
flow = partial(_flow, on_failure=[log_flow_failure])
文档
- State change hooksR G
11/13/2025, 8:20 AMMarvin
11/13/2025, 8:20 AMMarvin
11/13/2025, 8:21 AMon_failure 需要是列表(可包含多个钩子),而不是单个函数;
- 需要兼容两种用法:`@task` 和 `@task(...)`。因此要写成一个“既可直接装饰函数、也可带参数”的包装器,并且把默认钩子和用户自定义的钩子合并。
给你一份可直接用的“带默认 on_failure 的 task 包装器”,同时预设你提到的其它默认项:
from prefect import task as _prefect_task, get_run_logger
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
from prefect.cache_policies import NO_CACHE
def log_task_failure(task, task_run, state):
# 统一的失败钩子:打印原始异常 + traceback
exc = get_state_exception(state)
logger = get_run_logger()
logger.error(
f"[TASK FAILED] task={task.name} run_id={task_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
def task(__fn=None, **kwargs):
# 统一默认项
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
# 合并 on_failure:默认钩子 + 用户传入(可为单个或列表)
default_hooks = [log_task_failure]
user_hooks = kwargs.pop("on_failure", None)
if user_hooks is None:
hooks = default_hooks
elif isinstance(user_hooks, (list, tuple)):
hooks = default_hooks + list(user_hooks)
else:
hooks = default_hooks + [user_hooks]
kwargs["on_failure"] = hooks
def _decorate(fn):
return _prefect_task(**kwargs)(fn)
# 支持 @task 和 @task(...)
return _decorate if __fn is None else _decorate(__fn)
用法
@task
def task_a(x):
...
@task(retries=3)
def task_b(x):
...
def extra_hook(task, task_run, state): ...
@task(on_failure=[extra_hook], timeout_seconds=30)
def task_c(x):
...
补充说明
- 你之前的 kwargs.setdefault("on_failure", log_task_failure) 会导致:
- on_failure 不是列表;
- 丢失“合并”能力(用户传 on_failure 时会完全覆盖默认钩子而不是追加)。
- NO_CACHE 需从 prefect.cache_policies 引入;如果你后续需要结果持久化或跨运行读取,请把 persist_result 和存储/序列化配置回开。
- 如果你还想对 Flow 也默认化 on_failure,可以用同样方式包装 `prefect.flow`。
文档
- State change hooks(on_failure 用法)R G
11/13/2025, 9:46 AMMarvin
11/13/2025, 9:46 AMMarvin
11/13/2025, 9:47 AM@task 返回的是一个 Task 对象,它不是描述符(不实现 `__get__`),所以作为类属性时不会自动绑定 `self`。写成 self.method.submit(...) 时,`self` 不会被自动作为第一个参数传入,导致行为异常。
给你三种可用方案,按简单程度排序:
方案 1(最简单):调用时用“类属性”而不是“实例属性”
- 定义依然可以写在类里,但调用时用 ClassName.method.submit(self, ...) 显式传入 `self`。
class Svc:
@task # 你的带默认 on_failure 的包装
def step(self, x):
...
svc = Svc()
# 不要用 svc.step.submit(x)
# 要这样:
Svc.step.submit(svc, x) # 显式把 self 作为第一个参数传进去
方案 2(推荐):做一个“对方法友好”的 task 包装器(实现描述符注入 self)
- 下面的包装器在实例访问时返回一个代理,自动把 self 作为第一个参数注入,支持 `.submit()`、直接调用等;并且带上你要的默认参数与默认 on_failure。
from prefect import task as _prefect_task, get_run_logger
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
from prefect.cache_policies import NO_CACHE
def log_task_failure(task, task_run, state):
exc = get_state_exception(state)
logger = get_run_logger()
logger.error(
f"[TASK FAILED] task={task.name} run_id={task_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
class _BoundTaskProxy:
def __init__(self, task_obj, instance):
self._task = task_obj
self._instance = instance
# 直接调用:self.step(...)
def __call__(self, *args, **kwargs):
return self._task(self._instance, *args, **kwargs)
# 并发提交:self.step.submit(...)
def submit(self, *args, **kwargs):
return self._task.submit(self._instance, *args, **kwargs)
# 其他属性/方法透传
def __getattr__(self, name):
return getattr(self._task, name)
class _MethodFriendlyDescriptor:
def __init__(self, task_obj):
self._task = task_obj
def __get__(self, obj, owner):
# 通过类访问:返回原始 Task(适合 Class.method.submit(self, ...) 用法)
if obj is None:
return self._task
# 通过实例访问:返回自动注入 self 的代理(支持 self.method.submit(...))
return _BoundTaskProxy(self._task, obj)
def task(__fn=None, **kwargs):
# 统一默认项
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
# 合并 on_failure:默认 + 用户
default_hooks = [log_task_failure]
user_hooks = kwargs.pop("on_failure", None)
if user_hooks is None:
hooks = default_hooks
elif isinstance(user_hooks, (list, tuple)):
hooks = default_hooks + list(user_hooks)
else:
hooks = default_hooks + [user_hooks]
kwargs["on_failure"] = hooks
def _decorate(fn):
t = _prefect_task(**kwargs)(fn)
return _MethodFriendlyDescriptor(t)
return _decorate if __fn is None else _decorate(__fn)
用法:
class Svc:
@task
def step(self, x):
return x * 2
svc = Svc()
# 现在两种写法都可以:
r1 = svc.step.submit(10) # 自动注入 self
r2 = Svc.step.submit(svc, 10) # 显式传 self
方案 3:改成静态方法/外部函数
- 不依赖实例状态时,用 @staticmethod 并把需要的上下文通过参数传入;或直接定义为模块级函数,最稳妥。
- 例如:
class Svc:
@staticmethod
@task
def step(x, ctx):
...
注意装饰器顺序要保证传给 @task 的是可调用对象;如有冲突,优先改为模块级函数更简单。
另外,你的代码里还有几处小问题需要修正
- 形参与变量名不一致:你写的是 `def task(___fn_=None, **_kwargs_)`,但内部用的是 `__fn`、`kwargs`、`prefect_task`。请保证一致(参考上面方案 2 的实现)。
- on_failure 必须是列表;用 setdefault("on_failure", [log_task_failure]) 或合并逻辑。
- NO_CACHE 请从 prefect.cache_policies 引入。
如果你愿意,我可以帮你把现有项目的包装器替换成“方案 2”,并给你一个最小可运行样例做验证。R G
11/13/2025, 11:32 AMMarvin
11/13/2025, 11:32 AMMarvin
11/13/2025, 11:32 AMfrom prefect import task as _prefect_task, get_run_logger
from prefect.cache_policies import NO_CACHE
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
# 统一失败钩子
def log_task_failure(task, task_run, state):
exc = get_state_exception(state)
logger = get_run_logger()
logger.error(
f"[TASK FAILED] task={task.name} run_id={task_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
# 如果要同步写你们的日志系统:
# import logger.log as mario_log
# mario_log.warning(f"act=task_run_failed task={task.name} exc={exc}\n{exception_traceback(exc)}")
class _BoundTaskProxy:
def __init__(self, task_obj, instance):
self._task = task_obj
self._instance = instance
# 直接调用:self.step(...)
def __call__(self, *args, **kwargs):
return self._task(self._instance, *args, **kwargs)
# 并发提交:self.step.submit(...)
def submit(self, *args, **kwargs):
return self._task.submit(self._instance, *args, **kwargs)
# 支持 with_options:返回新的代理,继续保持自动注入 self
def with_options(self, **opts):
new_task = self._task.with_options(**opts)
return _BoundTaskProxy(new_task, self._instance)
# 其他属性/方法透传
def __getattr__(self, name):
return getattr(self._task, name)
class _MethodFriendlyDescriptor:
def __init__(self, task_obj):
self._task = task_obj
def __get__(self, obj, owner):
# 通过类访问:返回原始 Task(适合 Class.method.submit(self, ...) 用法)
if obj is None:
return self._task
# 通过实例访问:返回自动注入 self 的代理(支持 self.method.submit(...))
return _BoundTaskProxy(self._task, obj)
def task(__fn=None, **kwargs):
# 统一默认项
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
# 合并 on_failure:默认钩子 + 用户传入(可单个或列表)
default_hooks = [log_task_failure]
user_hooks = kwargs.pop("on_failure", None)
if user_hooks is None:
hooks = default_hooks
elif isinstance(user_hooks, (list, tuple)):
hooks = default_hooks + list(user_hooks)
else:
hooks = default_hooks + [user_hooks]
kwargs["on_failure"] = hooks
def _decorate(fn):
t = _prefect_task(**kwargs)(fn)
return _MethodFriendlyDescriptor(t)
# 支持 @task 和 @task(...)
return _decorate if __fn is None else _decorate(__fn)
用法示例(类成员函数)
class Svc:
@task # 自动带默认 on_failure
def step(self, x):
if x == 2:
raise RuntimeError("blade_req_final_failed=x=2")
return x * 2
svc = Svc()
# 方式 1:实例上直接 submit(自动注入 self)
f1 = svc.step.submit(2)
# 方式 2:类上 submit(显式传 self)
f2 = Svc.step.submit(svc, 2)
# 需要额外 options 时
f3 = svc.step.with_options(retries=3).submit(3)
为什么你之前的版本没生效
- 变量名混用导致 on_failure 实际没有传入 Prefect 的 task 构造器,比如 kwargs 和 `_kwargs_`、`__fn` 和 `___fn_`、`state` 和 _state_ 混搭。
- 描述符/代理里大量使用了未定义的变量名(如 `task_obj`/`name`/`args` 等),使得包装失败或行为异常。
- 修正后,on_failure 是真正传给了底层 Task,对类成员方法和普通函数都生效。
注意事项
- Prefect 的 Task 不是描述符,直接用 @prefect.task 装饰类方法时,实例上调用不会自动注入 self;这就是为什么我们加了 _MethodFriendlyDescriptor + _BoundTaskProxy。
- 如果你不需要实例上直接调用,也可以不用描述符,直接用类名调用:ClassName.method.submit(self, ...),更简单。
参考
- State change hooksR G
11/13/2025, 11:36 AMMarvin
11/13/2025, 11:36 AMMarvin
11/13/2025, 11:37 AMfrom prefect import task as _prefect_task, get_run_logger
from prefect.cache_policies import NO_CACHE
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
# 默认失败钩子
def log_task_failure(task, task_run, state):
exc = get_state_exception(state)
logger = get_run_logger()
logger.error(
f"[TASK FAILED] task={task.name} run_id={task_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
# 如需你们的日志:import logger.log as mario_log; mario_log.warning(...)
class _BoundTaskProxy:
def __init__(self, task_obj, instance):
self._task = task_obj
self._instance = instance
def __call__(self, *args, **kwargs):
return self._task(self._instance, *args, **kwargs)
def submit(self, *args, **kwargs):
return self._task.submit(self._instance, *args, **kwargs)
def with_options(self, **opts):
new_task = self._task.with_options(**opts)
return _BoundTaskProxy(new_task, self._instance)
def __getattr__(self, name):
return getattr(self._task, name)
class _MethodFriendlyDescriptor:
def __init__(self, task_obj):
self._task = task_obj
def __get__(self, obj, owner):
if obj is None:
return self._task # 类访问:BladesTask.req
return _BoundTaskProxy(self._task, obj) # 实例访问:blades_client.req
def task(__fn=None, **kwargs):
# 统一默认项
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
# 合并默认 + 用户 on_failure
default_hooks = [log_task_failure]
user_hooks = kwargs.pop("on_failure", None)
if user_hooks is None:
hooks = default_hooks
elif isinstance(user_hooks, (list, tuple)):
hooks = default_hooks + list(user_hooks)
else:
hooks = default_hooks + [user_hooks]
kwargs["on_failure"] = hooks
def _decorate(fn):
t = _prefect_task(**kwargs)(fn)
return _MethodFriendlyDescriptor(t)
return _decorate if __fn is None else _decorate(__fn)
使用示例(符合你的写法)
from my_prefect_wrappers import task
class BladesTask:
@task # 自动带 on_failure=[log_task_failure]
def req(self, x):
if x == "boom":
raise RuntimeError(f"blade_req_final_failed=x={x}")
return {"ok": x}
blades_client = BladesTask()
# 你的用法:先取成员,再 submit
blades_req = blades_client.req # 这是 _BoundTaskProxy
fut1 = blades_req.submit(123)
print(fut1.result())
# 或直接一行
fut2 = blades_client.req.submit("boom") # 失败时会触发 log_task_failure
fut2.result() # 这里会抛 UpstreamTaskError;根因已在钩子里打印过
常见“用起来不行”的原因与排查
- 误用了 Prefect 原生 task:请确保全项目都从你的包装模块导入 `task`,而不是 `from prefect import task`。
- 实例属性遮蔽了描述符:不要在 __init__ 里写 `self.req = ...`,这会覆盖类属性上的描述符,导致 blades_client.req 不再是代理对象。
- 多个装饰器顺序不当:要保证传给我们 task(...) 的是原始函数对象。若还有其它装饰器,通常让 @task 靠近函数定义最外层更稳妥;有冲突建议改为模块级函数。
- 取用方式写成了调用:需要 `blades_client.req.submit(...)`,不是 `blades_client.req().submit(...)`。
- 想确认是否生效:打印类型
print(type(BladesTask.req), type(blades_client.req))
# 应该显示 _MethodFriendlyDescriptor / _BoundTaskProxy(或内部包装类型)
如果你把这段放进项目后依旧不行,麻烦贴一下具体的报错信息(TypeError/AttributeError 等)和你类中 req 的定义(是否在 init 里重名覆盖),我再帮你对症修正。R G
11/13/2025, 11:41 AMMarvin
11/13/2025, 11:41 AMMarvin
11/13/2025, 11:43 AMfrom prefect import task as _prefect_task, get_run_logger
from prefect.cache_policies import NO_CACHE
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
# 默认失败钩子
def log_task_failure(task, task_run, state):
exc = get_state_exception(state)
logger = get_run_logger()
logger.error(
f"[TASK FAILED] task={task.name} run_id={task_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
# 如需你们的日志系统:
# import logger.log as mario_log
# mario_log.warning(f"act=task_run_failed task={task.name} exc={exc}\n{exception_traceback(exc)}")
class _BoundTaskProxy:
def __init__(self, task_obj, instance):
self._task = task_obj
self._instance = instance
def __call__(self, *args, **kwargs):
return self._task(self._instance, *args, **kwargs)
def submit(self, *args, **kwargs):
return self._task.submit(self._instance, *args, **kwargs)
def with_options(self, **opts):
# 合并默认 on_failure + 用户 on_failure,避免覆盖默认钩子
user_hooks = opts.pop("on_failure", None)
if user_hooks is not None:
base_default = getattr(self._task, "_default_on_failure_hooks", [])
if isinstance(user_hooks, (list, tuple)):
merged = list(base_default) + list(user_hooks)
else:
merged = list(base_default) + [user_hooks]
opts["on_failure"] = merged
new_task = self._task.with_options(**opts)
return _BoundTaskProxy(new_task, self._instance)
def __getattr__(self, name):
return getattr(self._task, name)
class _MethodFriendlyDescriptor:
def __init__(self, task_obj):
self._task = task_obj
def __get__(self, obj, owner):
if obj is None:
return self._task
return _BoundTaskProxy(self._task, obj)
def task(__fn=None, **kwargs):
# 默认项
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
# 合并默认 + 用户 on_failure
default_hooks = [log_task_failure]
user_hooks = kwargs.pop("on_failure", None)
if user_hooks is None:
hooks = default_hooks
elif isinstance(user_hooks, (list, tuple)):
hooks = default_hooks + list(user_hooks)
else:
hooks = default_hooks + [user_hooks]
kwargs["on_failure"] = hooks
def _decorate(fn):
t = _prefect_task(**kwargs)(fn)
# 标记默认钩子,供实例 with_options 合并使用
t._default_on_failure_hooks = default_hooks
return _MethodFriendlyDescriptor(t)
return _decorate if __fn is None else _decorate(__fn)
使用与注意
- 你的用法是正确的:
- blades_client = BladesTask()
- blades_req = blades_client.req
- blades_req.submit(xxx)
- 如果需要额外 options,请优先用“实例上的 with_options”,这样会保留默认 on_failure 并追加你自己的:
- blades_client.req.with_options(retries=3, on_failure=[extra_hook]).submit(xxx)
- 避免在类的 init 里写 self.req = ...(会遮蔽描述符)
- 如果你通过“类”去调用并传 on_failure(如 BladesTask.req.with_options(on_failure=[...])),Prefect 的原生 with_options 会直接覆盖现有 hooks。此时默认钩子会被替换。建议在“实例上”调用 with_options,或在类级别也做一个类似的合并包装(需要再加一层代理)。
快速自检
- 失败可复现的最小示例:
- 在方法里 raise RuntimeError("boom")
- 在 flow 里 fut = blades_client.req.submit("x"); fut.result()
- 期望看到 log_task_failure 打印的错误与 traceback(Prefect UI 的任务日志里也能看到)
文档参考
- State change hooksR G
11/13/2025, 11:45 AMMarvin
11/13/2025, 11:45 AMMarvin
11/13/2025, 11:45 AMblades_client = BladesTask()
blades_req = blades_client.req # 这是代理对象
blades_req.submit(xxx)
这就是你要的形式,on_failure 默认钩子也会生效。
方案 2:保持原生 Prefect @task,写一个轻量“绑定器”把类方法 Task 绑定到实例
- 适用于 Svc.step 是 <class 'prefect.tasks.Task'> 的场景。
先写个通用绑定器:
def bind_task(task_obj, instance):
class _BoundTask:
def __init__(self, t, inst):
self._task = t
self._inst = inst
def submit(self, *args, **kwargs):
return self._task.submit(self._inst, *args, **kwargs)
def __call__(self, *args, **kwargs):
return self._task(self._inst, *args, **kwargs)
def with_options(self, **opts):
new_t = self._task.with_options(**opts)
return bind_task(new_t, self._inst)
def __getattr__(self, name):
return getattr(self._task, name)
return _BoundTask(task_obj, instance)
然后这样用:
class Svc:
@task # 这里可以是你带默认 on_failure 的包装,或原生 prefect.task
def step(self, x):
...
svc = Svc()
blades_req = bind_task(Svc.step, svc) # 预先绑定 self
blades_req.submit(xxx)
可选:做成类里的属性,使用更顺手
class Svc:
@task
def step(self, x):
...
@property
def blades_req(self):
return bind_task(Svc.step, self)
svc = Svc()
svc.blades_req.submit(xxx)
说明
- 绑定器不会改变 Task 的配置,on_failure 钩子沿用 Task 本身的设置;如果你用的是“默认带 on_failure 的 @task 包装”,这里同样有效。
- 注意不要在 init 里给同名属性赋值(例如 self.step = ... 或 self.blades_req = ...),否则会遮蔽类属性/属性方法,导致绑定逻辑失效。R G
11/13/2025, 11:51 AMMarvin
11/13/2025, 11:51 AMMarvin
11/13/2025, 11:52 AMstate 和 `task`,这会在钩子执行时抛 `NameError`,Prefect 会记录“运行 hook 出错”,你看不到预期的日志。
- 修正成使用传入的参数名即可。
2) @task 传参名写错
- 你用了 `_name_`、`_task_run_name_`,Prefect 接受的是 `name`、`task_run_name`。错误的参数名会导致装饰器报错或忽略配置。
- 同时 task_run_name 里的占位符必须与函数真实参数名一致。你写了 `{feat_id}`、`{cos_name}`,但方法参数是 `_feat_id_`、`_cos_name_`,这会导致格式化失败或取不到值。
3) 可选:构造器里变量名不一致
- __init__(..., _bns_=...) 里使用了 `bns`,而不是 `_bns_`,这会在首次实例化时报 `NameError`。
给你一份“可直接用”的最小可运行版本(包含修正 + 方法友好版 @task 包装器 + 默认 on_failure)
先放到公共模块里(比如 workflow/utils/task.py):
from prefect import task as _prefect_task, get_run_logger
from prefect.cache_policies import NO_CACHE
from prefect.states import get_state_exception
from prefect.exceptions import exception_traceback
# 统一失败钩子:确保参数名和内部引用一致
def log_task_failure(task, task_run, state):
exc = get_state_exception(state)
logger = get_run_logger()
logger.error(
f"[TASK FAILED] task={task.name} run_id={task_run.id} "
f"exc={exc}\n{exception_traceback(exc)}"
)
class _BoundTaskProxy:
def __init__(self, task_obj, instance):
self._task = task_obj
self._instance = instance
def __call__(self, *args, **kwargs):
return self._task(self._instance, *args, **kwargs)
def submit(self, *args, **kwargs):
return self._task.submit(self._instance, *args, **kwargs)
def with_options(self, **opts):
# 合并默认 on_failure(避免被覆盖)
user_hooks = opts.pop("on_failure", None)
if user_hooks is not None:
base_default = getattr(self._task, "_default_on_failure_hooks", [])
if isinstance(user_hooks, (list, tuple)):
merged = list(base_default) + list(user_hooks)
else:
merged = list(base_default) + [user_hooks]
opts["on_failure"] = merged
new_task = self._task.with_options(**opts)
return _BoundTaskProxy(new_task, self._instance)
def __getattr__(self, name):
return getattr(self._task, name)
class _MethodFriendlyDescriptor:
def __init__(self, task_obj):
self._task = task_obj
def __get__(self, obj, owner):
if obj is None:
return self._task
return _BoundTaskProxy(self._task, obj)
def task(__fn=None, **kwargs):
# 默认项
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
# 合并默认 + 用户 on_failure
default_hooks = [log_task_failure]
user_hooks = kwargs.pop("on_failure", None)
if user_hooks is None:
hooks = default_hooks
elif isinstance(user_hooks, (list, tuple)):
hooks = default_hooks + list(user_hooks)
else:
hooks = default_hooks + [user_hooks]
kwargs["on_failure"] = hooks
def _decorate(fn):
t = _prefect_task(**kwargs)(fn)
t._default_on_failure_hooks = default_hooks
return _MethodFriendlyDescriptor(t)
return _decorate if __fn is None else _decorate(__fn)
然后你的类这样写(注意修正 name/task_run_name,以及参数占位符要与真实参数名一致;建议去掉下划线前后缀,避免踩坑):
```
from workflow.utils.task import task # 用你包装过的 task
class BladesTask:
def __init__(self, bns=blades_bns_online):
global MIAOBI_BLADES_BNS_CLIENT
global MIAOBI_BLADES_BNS_CLIENT_LOCK
if MIAOBI_BLADES_BNS_CLIENT is None:
with MIAOBI_BLADES_BNS_CLIENT_LOCK:
if MIAOBI_BLADES_BNS_CLIENT is None:
MIAOBI_BLADES_BNS_CLIENT = bns_client.BnsClient(bns)
self._blades_bns_client = MIAOBI_BLADES_BNS_CLIENT
@task(name="blades_req", task_run_name="blades-{feat_id}-{cos_name}")
async def req(
self,
feat_id,
user_id,
token,
group_id=0,Marvin
11/13/2025, 11:52 AM使用方式(你想要的 blades_req.submit):
blades_client = BladesTask()
blades_req = blades_client.req # 这是 _BoundTaskProxy
fut = blades_req.submit(
feat_id=1, user_id=2, token="xxx", cos_name="abc"
)
fut.result() # 这里会抛 UpstreamTaskError;根因日志会由 on_failure 打印
```
快速排查清单
- 确保 import 用的是你的包装器:`from workflow.utils.task import task`,不要混用 from prefect import task
- 不要在 init 里给 self.req 赋值,避免遮蔽类属性上的描述符
- @task 的参数名必须是 `name`、`task_run_name`,占位符与函数参数名一致
- 你的 on_failure 钩子里引用的变量名要和形参一致(上面的 log_task_failure 已修正)
- 有重试时,on_failure 只在“最终失败”时触发
如果按上面改完仍“没看到失败钩子日志”,请把 Prefect 任务日志里的提示贴一下(Prefect 在运行 hook 前后会打印内置日志,如 “Running hook ... in response to entering state Failed”/“An error was encountered while running hook ...”),我再帮你精确定位。R G
11/14/2025, 3:23 AMMarvin
11/14/2025, 3:23 AMMarvin
11/14/2025, 3:24 AM@task 装饰在类方法上,并且方法里依赖 `self._blades_bns_client`。要改成“普通函数 task”,核心是把状态(`bns_client` 实例)从类里移出去:要么作为参数传入,要么做成模块级单例。下面给你一个等价的函数版写法(以 Prefect 3.x 为例)。
要点
- 不再使用类;`@task` 直接装饰一个 async 函数。
- 用模块级单例获取并复用 `bns_client`(线程安全)。
- 修正你原代码里无效的参数名(例如 `@task(name=...)`→`@task(name=...)`;`httpx.AsyncClient(timeout=...)`→`httpx.AsyncClient(timeout=...)` 等)。
- 依赖注入:`fs`、`mario_log` 这类外部对象做成可选参数,便于测试和复用。
示例实现(功能与原方法一致,非必要的外部依赖做了可选注入)
```
import asyncio
import base64
import json
import threading
import traceback
import httpx
from prefect import task, get_run_logger
# 你原来的 BNS 地址(去掉了 <> 包裹)
BLADES_BNS_ONLINE = "bns://group.smartbns-from_product=default%group.blades-gateway-online.ps-wcb.all"
# 全局单例 + 线程锁,替代类上的成员
_MIAOBI_BLADES_BNS_CLIENT = None
_MIAOBI_BLADES_BNS_CLIENT_LOCK = threading.Lock()
def get_blades_bns_client(bns: str = BLADES_BNS_ONLINE):
global _MIAOBI_BLADES_BNS_CLIENT
if _MIAOBI_BLADES_BNS_CLIENT is None:
with _MIAOBI_BLADES_BNS_CLIENT_LOCK:
if _MIAOBI_BLADES_BNS_CLIENT is None:
# 这里沿用你原本的 bns_client 模块
_MIAOBI_BLADES_BNS_CLIENT = bns_client.BnsClient(bns)
return _MIAOBI_BLADES_BNS_CLIENT
@task(name="blades_req", task_run_name="blades-{feat_id}-{cos_name}")
async def blades_req(
feat_id,
user_id,
token,
group_id: int = 0,
video_url: str | None = None,
text_utf8: str | dict | None = None,
feat_args: dict | None = None,
ip_port: str | None = None,
retry_max: int = 3,
timeout_s: int = 60,
cos_name: str = "req",
log_id: int = 0,
output_format: str = "auto",
save_file: bool = False,
# 可选依赖注入
fs=None,
mario_log=None,
bns: str = BLADES_BNS_ONLINE,
):
"""
将原来的类方法重写为普通函数 task;bns client 通过模块级单例获取
"""
upload_log = get_run_logger()
# 可选:如果你有 get_local_ip(),可以保留;否则可以注释或自己实现
try:
upload_log.info(f"local_ip={get_local_ip()}")
except Exception:
pass
blade_info = f"feat_id={feat_id}, user_id={user_id}, token={token}, cos_name={cos_name}, log_id={log_id}"
upload_log.info(blade_info)
text_utf8_value = text_utf8
if isinstance(text_utf8_value, dict):
text_utf8_value = json.dumps(text_utf8_value, ensure_ascii=False)
# 组装请求
req_obj = {
"user": {"user_id": user_id, "token": token, "group_id": group_id},
"log_id": log_id,
"feat_args": [{"feat_id": feat_id, "feat_args": feat_args}],
"input_data": {"video_url": video_url, "text_utf8": text_utf8_value},
}
req_str = json.dumps(req_obj).encode("utf-8")
if save_file and fs is not None:
save_obj = text_utf8 if text_utf8 is not None else req_str
# 按你们自有 fs.save 的签名来,这里给出一种常见写法
fs.save(query=log_id, file_name=f"blades-{feat_id}-{cos_name}_input", text=save_obj)
# 获取 bns client(单例)
blades_bns_client = get_blades_bns_client(bns=bns)
result = None
msg = ""
async with httpx.AsyncClient(timeout=timeout_s) as http_client:
for retry_time in range(retry_max):
await asyncio.sleep(0.5 + 2 * retry_time)
try:
# 如果没有固定 ip:port,就从 bns 拉一个
if ip_port is None:
ip, port = blades_bns_client.get_a_host()
ip_port = f"{ip}:{port}"
blades_url = f"http://{ip_port}/BladesService/feat_calc"
upload_log.info(f"blades_ip={ip_port}, retry_time={retry_time}")
if mario_log:
mario_log.notice(f"act=blade_req_post, {blade_info}, retry_time={retry_time}")
# 异步请求
response = await http_client.post(
blades_url,
content=req_str,
headers={"Content-Type": "application/json"},
)
resp_json = response.json()
feat_status = resp_json["feat_res"][0]["feat_detail_status"]
# 非 0 视为失败Marvin
11/14/2025, 3:24 AM如何在 Flow 里调用
from prefect import flow
@flow
async def my_flow():
res = await blades_req(
feat_id=123,
user_id="u1",
token="xxx",
text_utf8={"q": "hello"},
# fs=..., mario_log=..., 视需要注入
)
return res
说明
- bns client 的复用以前靠类成员,现在通过 `get_blades_bns_client` 的模块级单例实现。
- 你原代码里像 `_name_`、`_timeout_`、`_headers_`、`_content_` 这种参数名是无效的,已改回标准参数名,否则 Prefect/httpx 都不会识别。
- 外部依赖(如 `fs.save`、`mario_log`、`get_local_ip`)如果不是全局可用,建议作为参数注入或在函数内做 try/except 保护。
如果你不想用全局单例,也可以把 `bns_client` 作为参数传入
@task
async def blades_req(..., bns_client):
...
if ip_port is None:
ip, port = bns_client.get_a_host()
```
需要我帮你把上面的代码对接你现有的 `fs.save`、`mario_log`、`get_local_ip` 的具体签名吗?或者你想保留单例还是改为参数注入?提示下你的偏好我可以给到最终版。
参考:Prefect Tasks(3.x 文档)R G
11/14/2025, 3:43 AMMarvin
11/14/2025, 3:43 AMMarvin
11/14/2025, 3:44 AM_decorate(_fn_) 里返回的是 `prefect_task(**kwargs)(fn)`,但这里的 fn 未定义,应该用 _fn_
- 最后一行 `return _decorate if __fn is None else _decorate(__fn)`,形参是 `___fn_`,用了不存在的 __fn
- log_task_failure 的参数是 `_task_/_task_run_/_state_`,但函数体里用了 task/state 变量名,导致 NameError
- 这些问题会让包装器失效,看起来像是“async 不支持”
按 Prefect 3.x 正确写法,下面这一版就可以完美支持同步/异步 task(`@task` 或 @task(...) 两种用法都支持):
from typing import Optional, Callable
from prefect import task as prefect_task
from prefect.cache_policies import NO_CACHE
from prefect.states import get_state_exception, State
from prefect.exceptions import exception_traceback
def log_task_failure(task, task_run, state: State):
"""
失败钩子:只会在最终失败时触发
"""
exc = get_state_exception(state)
tb = ""
try:
if exc:
tb = exception_traceback(exc)
except Exception:
pass
# 你自己的日志
try:
mario_log.warning(f"act=task_run_failed task={task.name} exc={exc}\n{tb}")
except NameError:
# 没有 mario_log 也不阻塞
pass
def task(__fn: Optional[Callable] = None, **kwargs):
"""
你的统一包装:支持 @task 和 @task(...)
同时支持 sync/async 函数
"""
# 统一默认值
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
kwargs.setdefault("on_failure", [log_task_failure])
def _decorate(fn: Callable):
# 直接委托给 Prefect 的官方装饰器(它本身支持 async)
return prefect_task(**kwargs)(fn)
# 支持两种用法:@task / @task(...)
return _decorate(__fn) if callable(__fn) else _decorate
使用示例
@task
async def my_async_task(x: int) -> int:
await asyncio.sleep(0.1)
return x + 1
@task(name="plus-two", retries=2)
def my_sync_task(x: int) -> int:
return x + 2
在 Flow 中调用(Prefect 3.x)
- 异步 Flow:直接 await
- 同步 Flow:直接调用即可(或用 .submit() 获得 future)
from prefect import flow
@flow
async def af():
y = await my_async_task(1)
z = my_sync_task(2) # 同步 task 直接返回结果
return y + z
@flow
def sf():
z = my_sync_task(2)
y = my_async_task(1) # 在同步 flow 中,async task 会被引擎调度执行,直接返回结果
return y + z
几点说明
- Prefect 的 task 装饰器天然支持 async 函数;只要你的包装器把函数原样交给 prefect.task(**kwargs)(fn) 就没问题。
- on_failure 的 hook 签名是 `(task, task_run, state)`;你可以同步实现(上面那样就行)。
- 你原来的问题主要是变量名不一致导致装饰器失效。
官方文档参考:Prefect 3.x Tasks 文档
如果你还想加默认的 `retries`、`retry_delay_seconds`、`task_run_name` 等,也可以在 kwargs.setdefault 里统一配置。需要我帮你按你们团队的默认策略再加几项吗?R G
11/14/2025, 3:48 AMMarvin
11/14/2025, 3:48 AMMarvin
11/14/2025, 3:49 AMtask 装饰器对异步函数天然支持。通常出现“好像不支持 async”的原因有这几类:
1) 没在 Flow 里运行
- 在 Flow 外调用 async task,本质上还是普通的异步函数,调用会返回 coroutine,需要 await 或 `asyncio.run`。
- 如果直接调不 `await`,就会看到 “coroutine was never awaited”。
2) 导入或命名遮蔽问题
- 需要确保你有 `from prefect import task as prefect_task`,否则你自定义的 task 会遮蔽 Prefect 的同名导入,最终递归或找不到真正的装饰器。
3) Hook 和依赖导入不完整
- get_state_exception 要从 prefect.states 导入;`exception_traceback` 要从 prefect.exceptions 导入。没导入会抛错,从而误以为 async 不支持。
给你一份完整可用的版本,支持 @task 和 @task(...),也同时支持 sync/async:
from typing import Callable, Optional
from prefect import task as prefect_task
from prefect.cache_policies import NO_CACHE
from prefect.states import State, get_state_exception
from prefect.exceptions import exception_traceback
def log_task_failure(task, task_run, state: State):
"""在最终失败时触发的 hook"""
exc = get_state_exception(state)
try:
tb = exception_traceback(exc) if exc else ""
except Exception:
tb = ""
try:
mario_log.warning(f"act=task_run_failed task={task.name} exc={exc}\n{tb}")
except NameError:
# 没有 mario_log 也不阻塞
pass
def task(__fn: Optional[Callable] = None, **kwargs):
"""自定义 task 装饰器:统一默认值 + 兼容 sync/async"""
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
kwargs.setdefault("on_failure", [log_task_failure])
def _decorate(fn: Callable):
return prefect_task(**kwargs)(fn)
# 支持 @task 和 @task(...)
return _decorate(__fn) if callable(__fn) else _decorate
使用方式示例
- 异步 Flow 里:需要 await
from prefect import flow
import asyncio
@task
async def my_async_task(x: int) -> int:
await asyncio.sleep(0.1)
return x + 1
@flow
async def aflow():
y = await my_async_task(1) # 注意 await
return y
- 同步 Flow 里:直接调用即可(引擎会调度 async 执行并返回结果)
from prefect import flow
@flow
def sflow():
y = my_async_task(1) # 不需要 await
return y
- Flow 外(纯 Python 环境):async 任务需要 await 或 asyncio.run
# 在交互式/脚本里这么用
import asyncio
asyncio.run(my_async_task(1))
如果你还是遇到“不支持 async”,请贴一下调用场景(在 Flow 里还是 Flow 外)和完整报错(尤其是是否有 “coroutine was never awaited”)。我可以据此进一步定位。文档参考:Prefect 3.x TasksR G
11/14/2025, 4:53 AMMarvin
11/14/2025, 4:53 AMMarvin
11/14/2025, 4:54 AMon_failure hook 是同步调用的(引擎没有 `await`),所以在 hook 里不能直接 `await`。
推荐两种安全修法(均为同步,兼容 async/sync task):
方案 A:直接使用 Prefect 提供的格式化函数(最简洁,零 await)
from prefect.states import format_exception
def log_task_failure(task, task_run, state):
mario_log.warning(
f"act=task_run_failed task={task.name} {format_exception(state)}"
)
方案 B:尽量拿到 Exception;拿不到就回退到字符串
import inspect
from prefect.states import get_state_exception, format_exception
from prefect.exceptions import exception_traceback
def log_task_failure(task, task_run, state):
exc = None
try:
exc = get_state_exception(state) # 有些环境下可能返回 coroutine
except Exception:
pass
if isinstance(exc, BaseException) and not inspect.iscoroutine(exc):
try:
tb = exception_traceback(exc)
except Exception:
tb = ""
msg = f"exc={exc}\n{tb}"
else:
# 没拿到 Exception(或拿到的是 coroutine),回退到字符串
msg = format_exception(state)
mario_log.warning(f"act=task_run_failed task={task.name} {msg}")
注意事项
- 不要把 coroutine 传给 `exception_traceback()`;它只接受真正的 BaseException 实例。
- on_failure hook 在 Prefect 引擎里是同步调用的,尽量避免在 hook 中做需要 await 的操作。
- 确保导入的是这些同步工具:
- from prefect.states import get_state_exception, format_exception
- from prefect.exceptions import exception_traceback
如果你确实更希望在 hook 内拿到“真实 Exception”,但当前 get_state_exception 返回的是 coroutine,可以改成纯字符串方案(方案 A),或者在 hook 外(比如在 task 内部)捕获并记录异常。
参考文档:
- Tasks(on_failure 钩子示例)
- States 与异常/格式化函数R G
11/14/2025, 5:05 AMMarvin
11/14/2025, 5:05 AMMarvin
11/14/2025, 5:05 AMtask_run 就是当前任务运行实例,你可以直接拿到运行名:
- 任务定义时设置的 `task_run_name`(包括用参数格式化后的最终名字),在运行时会被解析并存到 task_run.name
- task.name 是任务定义名(非运行名)
示例:
from prefect.states import format_exception
def log_task_failure(task, task_run, state):
run_name = getattr(task_run, "name", None) or "(unknown)"
task_name = getattr(task, "name", None) or "(unknown)"
msg = f"act=task_run_failed task={task_name} run_name={run_name} {format_exception(state)}"
try:
mario_log.warning(msg)
except NameError:
print(msg)
定义任务时设置运行名(动态参数会被引擎替换,hook 里拿到的是替换后的最终名字):
from prefect import task
@task(task_run_name="blades-{feat_id}-{cos_name}")
async def blades_req(feat_id, cos_name="req"):
...
参考文档:
- Prefect Tasks(3.x)R G
11/14/2025, 5:07 AMMarvin
11/14/2025, 5:07 AMMarvin
11/14/2025, 5:08 AMfrom prefect.states import format_exception
def log_task_failure(task, task_run, state):
run_name = getattr(task_run, "name", "(unknown)")
msg = f"act=task_run_failed task={task.name} run_name={run_name} {format_exception(state)}"
try:
mario_log.warning(msg)
except NameError:
print(msg)
如果你坚持拿到 Exception 再打印 traceback(兼容 async/sync)
- 某些环境下 get_state_exception(state) 可能返回 coroutine。在 hook(同步函数)里不能直接 await,最好回退到 `format_exception(state)`。下面给一个“尽力拿 Exception,否则回退”的写法:
import inspect
from prefect.states import get_state_exception, aget_state_exception, format_exception
from prefect.exceptions import exception_traceback
def log_task_failure(task, task_run, state):
run_name = getattr(task_run, "name", "(unknown)")
exc = None
try:
exc = get_state_exception(state)
except Exception:
exc = None
# 某些环境下会返回 coroutine;尝试桥接,否则回退
if inspect.iscoroutine(exc):
try:
import anyio
# 如果当前在工作线程里,这里能把协程交回事件循环执行
exc = anyio.from_thread.run(aget_state_exception, state)
except Exception:
exc = None
if isinstance(exc, BaseException):
try:
tb = exception_traceback(exc)
except Exception:
tb = ""
msg = f"exc={exc}\n{tb}"
else:
# 拿不到 Exception(或不在可桥接的上下文)时,用官方格式化,包含完整堆栈
msg = format_exception(state)
try:
mario_log.warning(f"act=task_run_failed task={task.name} run_name={run_name} {msg}")
except NameError:
print(f"act=task_run_failed task={task.name} run_name={run_name} {msg}")
要点
- on_failure hook 是同步函数;不要在里面直接 `await`。
- format_exception(state) 已经包含完整的异常/堆栈信息,适用于 async/sync 任务。
- 如果你确实要用 `exception_traceback`,务必先确认拿到的是 `BaseException`,而不是一个 coroutine;否则就回退到 `format_exception(state)`。
参考:
- Tasks(on_failure 钩子)
- States(异常与格式化)R G
11/14/2025, 5:11 AMMarvin
11/14/2025, 5:11 AMMarvin
11/14/2025, 5:13 AMprefect.exceptions.exception_traceback(exc) 生成完整堆栈(含行号)。困难在于某些环境下 get_state_exception(state) 返回的是 coroutine;on_failure 是同步调用,不能直接 await。可以用 anyio 的桥接在同步钩子里“间接 await”。
给你一份稳健实现(优先拿 Exception;拿不到就兜底用 `format_exception`):
import inspect
from prefect.states import get_state_exception, aget_state_exception, format_exception
from prefect.exceptions import exception_traceback
def _extract_exception_sync(state):
"""
尝试同步地从 State 中提取 BaseException。
如遇到 coroutine,使用 anyio.from_thread.run 桥接到当前事件循环。
"""
try:
exc = get_state_exception(state)
except Exception:
exc = None
# 有些环境这里会是 coroutine,需要桥接
if inspect.iscoroutine(exc):
try:
import anyio
exc = anyio.from_thread.run(aget_state_exception, state)
except Exception:
exc = None
return exc if isinstance(exc, BaseException) else None
def log_task_failure(task, task_run, state):
run_name = getattr(task_run, "name", "(unknown)")
exc = _extract_exception_sync(state)
if exc is not None:
# 完整 traceback(含行号/文件)
try:
tb = exception_traceback(exc)
except Exception:
tb = ""
msg = f"exc={exc}\n{tb}"
else:
# 兜底:拿不到 Exception 时至少打印格式化信息
# 如你的版本没有完整堆栈,这里可考虑升级 Prefect,或将 trace 放到任务内 try/except 里打日志
msg = format_exception(state)
try:
mario_log.warning(f"act=task_run_failed task={task.name} run_name={run_name} {msg}")
except NameError:
print(f"act=task_run_failed task={task.name} run_name={run_name} {msg}")
说明
- exception_traceback(exc) 会返回和 Python traceback 风格一致的多行文本,包含文件名与行号。
- anyio.from_thread.run(...) 能在同步上下文里“调度执行”异步函数,适合 Prefect 引擎这类 anyio 管理的环境。
- 若你的环境无法桥接(极少见),可在任务内部使用 try/except + traceback.format_exc() 主动记录堆栈,或兜底用 `format_exception(state)`。
相关文档:
- Tasks(on_failure 钩子): https://docs.prefect.io/latest/concepts/tasks/
- States(异常与格式化): https://docs.prefect.io/latest/concepts/states/R G
11/14/2025, 5:17 AMMarvin
11/14/2025, 5:17 AMMarvin
11/14/2025, 5:18 AMraise_state_exception(state) 把异常抛出来,再在钩子里捕获并格式化。
给你一版稳健实现(不依赖异步,兼容你当前版本):
import inspect
from prefect.states import get_state_exception, format_exception
try:
from prefect.states import raise_state_exception # 同步抛出状态里的异常
except ImportError:
raise_state_exception = None
from prefect.exceptions import exception_traceback
def log_task_failure(task, task_run, state):
run_name = getattr(task_run, "name", "(unknown)")
exc = None
# 尝试 1:直接拿异常(某些版本这里可能返回 coroutine,就跳过)
try:
e = get_state_exception(state)
if isinstance(e, BaseException) and not inspect.iscoroutine(e):
exc = e
except Exception:
pass
# 尝试 2:用同步 API 把异常抛出来再捕获(拿到真实 Exception 和完整堆栈)
if exc is None and callable(raise_state_exception):
try:
try:
raise_state_exception(state) # 这里会 raise 真实异常
except BaseException as e:
exc = e
except Exception:
pass
if exc is not None:
# 优先用 Prefect 的 traceback 格式化;失败则回退到 Python 标准库
try:
tb = exception_traceback(exc)
except Exception:
import traceback
tb = "".join(traceback.format_exception(type(exc), exc, exc.__traceback__))
msg = f"exc={exc}\n{tb}"
else:
# 实在拿不到 Exception(或者版本不支持),兜底
msg = format_exception(state)
try:
mario_log.warning(f"act=task_run_failed task={task.name} run_name={run_name} {msg}")
except NameError:
print(f"act=task_run_failed task={task.name} run_name={run_name} {msg}")
说明
- raise_state_exception(state) 是同步的;在钩子里调用它会抛出任务里真实的异常,你在钩子里立刻捕获,就能拿到包含行号/文件的完整堆栈。
- 仍保留 get_state_exception 的尝试;如果它在你的版本中返回了 coroutine,我们会跳过。
- 最后兜底用 `format_exception(state)`,以免某些极端情况完全没有输出。
建议
- 告知你的 Prefect 版本(`prefect version`),不同小版本的 states 工具函数存在差异。若方便,升级到较新的 3.x 后可用 aget_state_exception 等更清晰的 API。
- 文档参考:
- Tasks(on_failure 钩子)
- States(异常相关工具)R G
11/14/2025, 5:21 AMMarvin
11/14/2025, 5:21 AMMarvin
11/14/2025, 5:21 AMimport inspect
import traceback as _traceback
from prefect.states import get_state_exception, raise_state_exception, format_exception
from prefect.exceptions import exception_traceback
def log_task_failure(_task_, _task_run_, _state_):
run_name = getattr(_task_run_, "name", "(unknown)")
exc = None
# 尝试 1:直接取异常;部分环境可能返回 coroutine,就跳过
try:
e = get_state_exception(_state_)
if isinstance(e, BaseException) and not inspect.iscoroutine(e):
exc = e
except Exception:
pass
# 尝试 2:同步抛出再捕获,拿到真实异常和完整堆栈(适配 async/sync)
if exc is None:
try:
try:
raise_state_exception(_state_) # 这里会 raise 真实异常
except BaseException as e:
exc = e
except Exception:
exc = None
if exc is not None:
# 优先 Prefect 的 traceback;失败则回退到 Python 标准库
try:
tb = exception_traceback(exc)
except Exception:
tb = "".join(_traceback.format_exception(type(exc), exc, exc.__traceback__))
msg = f"task={_task_.name} run_name={run_name} exc={exc}\n{tb}"
else:
# 实在拿不到 Exception(极少见),兜底
msg = f"task={_task_.name} run_name={run_name} {format_exception(_state_)}"
try:
mario_log.warning(msg)
except NameError:
print(msg)
说明
- raise_state_exception(state) 是同步 API,会把任务里实际抛出的异常重新 raise;你在钩子里立刻捕获,就能拿到完整 traceback(含文件/行号)。这对 async 任务同样有效。
- get_state_exception 在某些打包环境里可能返回 coroutine,此时不要直接传给 exception_traceback;上面代码会自动避开这种情况。
- 你也可以简化为只用 raise_state_exception 路径(省去 get_state_exception 的尝试)。
相关文档
- Tasks(on_failure 钩子): https://docs.prefect.io/latest/concepts/tasks/
- States(异常工具): https://docs.prefect.io/latest/concepts/states/R G
11/14/2025, 5:28 AMMarvin
11/14/2025, 5:28 AMMarvin
11/14/2025, 5:29 AMraise_state_exception(state) 在你这个版本/上下文里返回了一个 coroutine,但 on_failure 钩子是同步调用的,Prefect 不会帮你 `await`。解决办法:在钩子里用 anyio 的线程桥接把这个 coroutine 跑起来,再捕获真实异常并格式化 traceback(含文件/行号)。
给你一份稳健可用的实现(Prefect 3.4.12,适配 async/sync 任务,钩子保持同步):
import inspect
import traceback as _traceback
from prefect.states import get_state_exception, raise_state_exception
try:
# 如果你的版本提供了异步版,优先用它
from prefect.states import araise_state_exception
except Exception:
araise_state_exception = None
try:
from prefect.exceptions import exception_traceback
except Exception:
exception_traceback = None
def _format_exc_with_tb(exc: BaseException) -> str:
# 优先用 Prefect 的 traceback 格式化;没有就用标准库
if exception_traceback:
try:
return exception_traceback(exc)
except Exception:
pass
return "".join(_traceback.format_exception(type(exc), exc, exc.__traceback__))
def log_task_failure(_task_, _task_run_, _state_):
run_name = getattr(_task_run_, "name", "(unknown)")
exc = None
# 尝试 1:直接拿异常;如果拿到的是 coroutine 就跳过
try:
e = get_state_exception(_state_)
if isinstance(e, BaseException) and not inspect.iscoroutine(e):
exc = e
except Exception:
pass
# 尝试 2:让状态把异常“抛”出来,再在钩子里捕获(需要 anyio 桥接)
if exc is None:
try:
import anyio
if araise_state_exception is not None:
# 异步抛异常:在同步钩子里通过 anyio.from_thread.run 执行
try:
anyio.from_thread.run(araise_state_exception, _state_)
except BaseException as e:
exc = e
else:
# 同步 API 可能返回 coroutine;判断并桥接执行
res = None
try:
res = raise_state_exception(_state_) # 可能直接 raise,也可能返回 coroutine
except BaseException as e:
exc = e
if exc is None and inspect.iscoroutine(res):
async def _runner():
# 等待并让异常在这里抛出
await res
try:
anyio.from_thread.run(_runner)
except BaseException as e:
exc = e
except Exception:
# anyio 不可用或其他非关键失败,忽略
pass
if exc is not None:
tb_text = _format_exc_with_tb(exc)
msg = f"task={_task_.name} run_name={run_name} exc={exc}\n{tb_text}"
else:
# 实在拿不到 Exception(极少见),打印最基本的信息
msg = f"task={_task_.name} run_name={run_name} (no exception object available)"
try:
mario_log.warning(msg)
except NameError:
print(msg)
要点说明
- 不要在钩子里直接调用 raise_state_exception(state) 后不处理返回值;它可能是 coroutine。
- 通过 anyio.from_thread.run(...) 可以在同步钩子中安全运行异步函数或 coroutine。Prefect 引擎基于 anyio,运行环境里应当可用。
- 一旦拿到真实的 `BaseException`,用 prefect.exceptions.exception_traceback(exc) 或标准库 traceback.format_exception 输出完整堆栈(包含文件/行号)。
如果你更倾向于精简,也可以直接只用“尝试 2”那段(用 araise_state_exception 或桥接 `raise_state_exception`),去掉“尝试 1”。我保留“尝试 1”只是为了在某些环境下能更快拿到异常。
参考
- Prefect 3.x Tasks(hooks)
- Prefect 3.x States(异常相关)R G
11/14/2025, 7:48 AMquote annotation to skip resolution of inputs. 这是什么错误?Marvin
11/14/2025, 7:48 AMMarvin
11/14/2025, 7:53 AMfrom prefect import task, flow
from prefect.utilities.annotations import quote
@task
def do_something(server_record: quote):
# server_record 会被按原样传入,不会被 Prefect 解析/遍历
...
@flow
def my_flow(server_record):
do_something(server_record) # 这里直接传即可
- 在调用处对值做包装
from prefect.utilities.annotations import quote
do_something(server_record=quote(server_record))
其他容易踩坑的原因与对策
- 传入了协程/未完成的 Future:先 await/取结果再传
- 传入了不可序列化对象(DB 会话、HTTP 客户端、文件句柄等):用 quote 跳过解析,或改成可序列化的简单结构(dict/str 等)
- 复杂嵌套结构里只有一部分是“特殊对象”:只对那一部分做 quote 包裹即可
注意
- 被 quote 的参数会跳过 Prefect 的输入解析/映射逻辑,不能用于 map,也通常不会被结果持久化/可视化。
- 如果参数实际上应该是普通 JSON 可序列化数据,尽量转换为基础类型,而不是一味 quote。
参考文档
- 任务/参数注解概念:https://docs.prefect.io/latest/concepts/tasks/
- 流/参数注解概念:https://docs.prefect.io/latest/concepts/flows/#parameter-annotations
- 注解 API(quote/unmapped/allow_failure):https://docs.prefect.io/latest/api-ref/utilities/annotations/R G
11/17/2025, 10:15 AMquote annotation to skip resolution of inputs. 这个错误是什么原因?报错的代码是:record_data = baikan_tasks.preproc_upload_to_bos(
input_json, text_res, planning_tts_json, tts_res, ttv_res_list, video_path_list_url,
cover_res, card_res_list, aigc_video_res_list, voice_effect_res, video_effect_res)
server_record = upload_to_bos(record_data)。其中preproc_upload_to_bos和upload_to_bos都是task,预期record_data是一个json,并且我自己排查时,在preproc_upload_to_bos函数中打印了要输出record_data,看打印结果确实是jsonMarvin
11/17/2025, 10:15 AMMarvin
11/17/2025, 10:20 AMupload_to_bos 任务之前,Prefect 会递归“解析(resolve)”你传入的参数,把里面的 Prefect futures、可迭代容器等都遍历并转成最终值。如果在遍历 upload_data 这个参数时遇到 Prefect 不支持/无法遍历的类型,就会报 “Failed to resolve inputs in parameter 'upload_data'… consider using the quote annotation”。
结合你这段代码:
record_data = baikan_tasks.preproc_upload_to_bos(
input_json, text_res, planning_tts_json, tts_res, ttv_res_list, video_path_list_url,
cover_res, card_res_list, aigc_video_res_list, voice_effect_res, video_effect_res
)
server_record = upload_to_bos(record_data)
- record_data 是上游 task 的返回值(在 flow 里是个 PrefectFuture,Prefect 会在调用 upload_to_bos 之前解析为真正的结果)。
- 即使你在 preproc_upload_to_bos 里打印它看起来是个 “json/dict”,里面可能还嵌着 Prefect 不好处理的对象(例如 pandas DataFrame、numpy 数组、requests.Response、复杂自定义类、包含方法/文件句柄等)。Prefect 在“解析输入”时会递归进入 dict/list 等容器;若某个值不可遍历/不受支持,就会触发这个错误。
两种常见修复方式
1) 把 record_data 里的“复杂值”转成基础类型
- numpy 数组 -> tolist()
- pandas DataFrame -> to_dict(orient="records") 或 to_json()
- 自定义对象 -> 转成 `dict`/`str`
- 二进制内容 -> `bytes`/`str`(尽量避免让 Prefect递归进去,可见下条的 quote)
示例:
@task
def preproc_upload_to_bos(...):
arr = np.arange(3).tolist()
df = some_df.to_dict(orient="records")
return {"arr": arr, "df": df, "other": "ok"}
2) 用 quote 跳过对某些字段的解析
- 只给“有问题的字段”加 `quote`,不要把整个 `record_data`(或整个 future)都 quote;否则会把 PrefectFuture 原样传给下游任务,导致下游拿到的是 future 而不是实际值。
- 最好在上游 task 里对具体字段 quote,然后返回。
示例:
from prefect.utilities.annotations import quote
from prefect import task
@task
def preproc_upload_to_bos(...):
# 假设 df / binary_data / custom_obj 是解析时会出问题的值
record_data = {
"meta": {...},
"df": quote(df), # 跳过递归解析 df
"binary": quote(binary_data), # 比如 bytes
"custom": quote(custom_obj), # 自定义对象
}
return record_data
@task
def upload_to_bos(upload_data):
# 这里会收到完整的 dict,但 df/binary/custom 不会被 Prefect 递归展开
...
如何快速定位“哪个字段有问题”
- 在 preproc_upload_to_bos 里打印类型分布,帮助找出需要转换或 quote 的字段:
from prefect import get_run_logger
@task
def preproc_upload_to_bos(...):
record_data = {...}
logger = get_run_logger()
<http://logger.info|logger.info>({k: type(v).__name__ for k, v in record_data.items()})
return record_data
- 也可以逐步删减字段来二分定位元凶。
注意事项
- 不要把整个 record_data 的 PrefectFuture 用 quote(record_data) 包起来传给下游,否则下游会收到一个 future 对象而不是 dict。
- 如果你确实只想传字符串,可以在上游 `json.dumps(record_data)`,下游再 `json.loads`,这样也能绕过解析。
参考链接
- Annotations(含 `quote`):prefect.utilities.annotations
- 集合遍历/解析机制:prefect.utilities.collections
- 解析实现(源码,了解报错来源):resolve_inputs、visit_collection
如果你能贴一下 record_data 的结构或每个字段的类型(例如是否包含 pandas/numpy/自定义类/bytes),我可以更精准地建议是转型还是在哪些字段上使用 `quote`。R G
11/17/2025, 11:02 AMMarvin
11/17/2025, 11:02 AMMarvin
11/17/2025, 11:03 AMupload_to_bos(upload_data=record_data) 时,`record_data` 里面很可能还残留了“Prefect 不支持/无法遍历”的值(或残留的 future),导致 Prefect 在解析 upload_data 时失败。
重点解释
- Prefect 会在调用任务前,递归解析参数,把 futures 展开成最终结果,并遍历 dict/list 等容器里的每个元素。如果某个元素类型不受支持(或是未被解析的 future 容器),就会触发你看到的错误。
- 你贴出来的类型是传给上游任务 preproc_upload_to_bos 的输入,其中 ttv_res/card_res_list/aigc_video_res_list 是 PrefectFutureList,这没问题;Prefect 会把它们解析成真实结果再喂给任务函数。
- 但最终返回的 record_data 里,某个字段可能还是复杂对象(或残留的 future/future list),到 upload_to_bos 时解析失败。
建议排查与修复
1) 在 preproc_upload_to_bos 任务内部,打印 record_data 每个字段的叶子类型,确认是否还包含 PrefectFuture/PrefectFutureList 或其他复杂对象(如 pandas/numpy/自定义类/bytes/Response 等)
from prefect import get_run_logger
def _walk_types(x, path=""):
t = type(x).__name__
if isinstance(x, dict):
return {k: _walk_types(v, f"{path}.{k}") for k, v in x.items()}
elif isinstance(x, (list, tuple)):
return [ _walk_types(v, f"{path}[{i}]") for i, v in enumerate(x) ]
else:
return t
@task
def preproc_upload_to_bos(...):
record_data = {...} # 你构造的返回
logger = get_run_logger()
<http://logger.info|logger.info>(_walk_types(record_data))
return record_data
2) 把“问题字段”转成基础类型
- numpy 数组 -> tolist()
- pandas DataFrame -> to_dict(orient="records") 或 to_json()
- 自定义对象 -> `dict`/`str`
- bytes/binary -> `str`(或后面用 quote,见下)
3) 对个别字段用 quote 跳过深度解析(只 quote 有问题的字段)
from prefect.utilities.annotations import quote
@task
def preproc_upload_to_bos(...):
record_data = {
"ok_meta": {...},
"maybe_binary": quote(binary_data), # 避免 Prefect 递归进去
"custom_obj": quote(custom_obj), # 自定义对象也可 quote
}
return record_data
- 注意不要对整个 `record_data`(或对 record_data 这个 future)做 `quote`,否则下游会拿到 future 本体而不是数据。
- quote 参考:prefect.utilities.annotations
4) 如果仍不清楚问题字段在哪,可以把 record_data 逐步删减字段(二分法)来定位元凶;或把“可能复杂”的字段先 json.dumps 成字符串传递,下游再 `json.loads`。
回答你的问题:“这些类型有问题吗?”
- 作为 preproc_upload_to_bos 的输入,出现 PrefectFutureList 正常;Prefect 会解析它们。
- 真正需要关注的是 preproc_upload_to_bos 返回的 record_data 的内部各字段类型。只要全部是基础 Python/JSON 友好类型(dict/list/str/int/float/bool/None),或者对复杂字段使用了 `quote`,`upload_to_bos` 就不会在解析 upload_data 时报错。
如果你能贴出 record_data 的结构(每个字段的类型分布,尤其是那些不是 str/int/bool/list/dict 的值),我可以帮你精确指出哪些要转换或使用 `quote`。参考解析机制与遍历工具文档:
- 解析/遍历工具:prefect.utilities.collections
- 源码(了解报错来源):resolve_inputs、visit_collectionR G
11/17/2025, 11:05 AMMarvin
11/17/2025, 11:05 AMMarvin
11/17/2025, 11:06 AMupload_to_bos(upload_data=record_data) 解析 upload_data 时就会失败,报你看到的 “Failed to resolve inputs in parameter 'upload_data'…”。
修复方法(首选)
- 在 task 里用“参数本身”的名字,不要用全局变量。最好把参数名改成一致的、无下划线的名字,函数体里就用同名变量。
示例(修正后)
from prefect import task
@task
def preproc_upload_to_bos(
input_json, text_res, planning_tts_json, tts_res, ttv_res, video_path_list_url,
cover_res, card_res_list, aigc_video_res_list, voice_effect_res, video_effect_res
):
# 这里使用的都是“参数本身”,Prefect 在进入 task 前会先把 futures 解析为真实值
ret = {
"input_json": input_json,
"text_res": text_res,
"planning_tts_json": planning_tts_json,
"tts_res": tts_res,
"ttv_res": ttv_res, # 应该是普通 list/数据结构,不再是 PrefectFutureList
"video_path_list_url": video_path_list_url,
"cover_res": cover_res,
"card_res_list": card_res_list,
"aigc_video_res_list": aigc_video_res_list,
"voice_effect_res": voice_effect_res,
"video_effect_res": video_effect_res,
}
# 打印叶子类型,确认都已解析为基础类型
ret_type = {k: type(v).__name__ for k, v in ret.items()}
logger = get_run_logger()
<http://logger.info|logger.info>(ret_type)
# 如果 fs.save 需要字符串,记得序列化
# fs.save(_query_=input_json["task_id"], _file_name_="preproc_upload_to_bos", _text_=json.dumps(ret, ensure_ascii=False))
# fs.save(_query_=input_json["task_id"], _file_name_="preproc_upload_to_bos_type", _text_=json.dumps(ret_type, ensure_ascii=False))
return ret
如果修正后仍报错
- 说明 ret 的某些字段还是 Prefect 不好解析的复杂对象(比如 numpy/pandas/自定义类/bytes/HTTP Response 等)。两种办法:
1) 在上游 task 里把这些字段转成基础类型(list/dict/str);
2) 只对“问题字段”使用 quote 跳过解析(不要对整个 ret 或整个 future 使用 `quote`)。
示例(对个别字段用 quote)
from prefect.utilities.annotations import quote
@task
def preproc_upload_to_bos(...):
...
ret = {
"input_json": input_json,
"df": quote(df), # pandas DataFrame -> 跳过 Prefect 递归
"binary": quote(binary_bytes), # 二进制内容
"custom": quote(custom_obj), # 自定义对象
...
}
return ret
为什么你之前打印会看到 PrefectFutureList?
- 正是因为用了全局变量名(无下划线),绕开了 Prefect 的“参数解析流程”,导致你在 task 里看到的不是解析后的值,而是原始的 futures 容器。
参考
- 注解(含 `quote`):prefect.utilities.annotations
- 集合遍历/解析机制:prefect.utilities.collections
如果你把上面的参数名修正后再跑一次,把 ret_type 的日志贴出来(用“参数本身”的名字打印),我可以再帮你确认是否还需要对某些字段转型或 `quote`。R G
11/17/2025, 11:07 AMMarvin
11/17/2025, 11:08 AMMarvin
11/17/2025, 11:13 AMupload_to_bos(upload_data=record_data) 解析 upload_data 时失败的根因。
这通常有两种原因:
- 该 task 并没有在一个 @flow 运行上下文中执行(比如直接在普通 Python 代码里调用)。在这种情况下,参数不会被 Prefect 解析。
- 上游把一个 PrefectFutureList 直接当值往下传,且解析阶段没有被展开(例如某些场景/封装导致 Prefect 没有递归进去)。
快速修复方案(任选其一或组合):
1) 确保是在 flow 里调用
示例:
from prefect import flow
@flow
def main():
record_data = preproc_upload_to_bos(...所有上游结果...)
server_record = upload_to_bos(record_data)
if __name__ == "__main__":
main()
如果你本来就在 flow 里,请看第 2 点。
2) 在 task 开头显式把 Future 列表解析成结果列表
from prefect.futures import resolve_futures_to_results
from prefect import task
@task
def preproc_upload_to_bos(...):
ttv_res = resolve_futures_to_results(ttv_res)
card_res_list = resolve_futures_to_results(card_res_list)
aigc_video_res_list = resolve_futures_to_results(aigc_video_res_list)
ret = {
"input_json": input_json,
"text_res": text_res,
"planning_tts_json": planning_tts_json,
"tts_res": tts_res,
"ttv_res": ttv_res, # 现在是普通 list
"video_path_list_url": video_path_list_url,
"cover_res": cover_res,
"card_res_list": card_res_list, # 普通 list
"aigc_video_res_list": aigc_video_res_list, # 普通 list
"voice_effect_res": voice_effect_res,
"video_effect_res": video_effect_res
}
return ret
参考 API:prefect.futures
3) 如果这些 list 的元素里还有复杂/不可序列化对象(numpy/pandas/自定义类/二进制等),请在上游转换为基础类型,或者只对“问题字段”使用 quote 跳过解析(不要对整个 ret 使用 `quote`):
from prefect.utilities.annotations import quote
ret = {
...
"binary": quote(binary_bytes),
"custom": quote(custom_obj),
}
quote 文档:prefect.utilities.annotations
补充:
- 你在 fs.save(..., text=ret) 传的是一个 dict,如果 fs.save 期望的是字符串,建议 json.dumps(ret, ensure_ascii=False) 再写入,避免额外问题。
- 如果方便,建议把这一段的 flow 包装代码也贴一下;确认在 @flow 上下文里调用,能帮助我们排除第 1 个原因。R G
11/18/2025, 1:05 AMMarvin
11/18/2025, 1:05 AMMarvin
11/18/2025, 1:09 AMlog 表:行数巨大、删除会产生大量 dead tuples、需要频繁 VACUUM;再加上没有良好索引或按时间清理就更慢。给你几个层次的优化方案,按“立刻可用 → 架构级提速”排序。
快速提速(尽量不改架构)
- 先按时间直接清理日志,再清理 flow_run。对于“保留 12 天”的策略,日志和 flow_run 的生存期一致,优先按 timestamp 清理日志就能避免你现在的两次 JOIN 删除。
- 给日志相关列建合适的索引,特别是时间列用 BRIN,ID 列用 BTree:
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_timestamp_brin ON log USING brin (timestamp) WITH (pages_per_range = 32);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_flow_run_id ON log (flow_run_id);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_log_task_run_id ON log (task_run_id);
- 用批量分段删除,避免长事务和大锁,占用更少的 autovacuum 预算。常用模式(循环执行直到影响行数为 0):
-- 每次删 100k,走 timestamp 范围扫描(需要前面的索引)
WITH batch AS (
SELECT id
FROM log
WHERE timestamp < now() - interval '12 days'
ORDER BY timestamp
LIMIT 100000
)
DELETE FROM log l
USING batch
WHERE l.id = batch.id;
- 把日志清理放在 flow_run 清理之前,这样你后续对 log 的两次 JOIN 删除基本都能省掉(或只剩少量近 12 天内的数据)。
- 调整 autovacuum 针对大表的参数(特别是 `log`),加快回收:
- 降低 `autovacuum_vacuum_scale_factor`(例如 0.01 或更低)
- 设置较小的 `autovacuum_vacuum_threshold`(例如 5000)
- 高峰期外安排 `VACUUM (VERBOSE, ANALYZE) log;`,必要时对 log 做 `REINDEX CONCURRENTLY`(删除量极大且索引膨胀时)
- 如果你可以接受更少的 UI 日志:把 Prefect 发送日志到 API 的功能关掉或降量,源头减流量
- 完全关闭(不进数据库):`PREFECT_LOGGING_TO_API_ENABLED=false`
- 降低日志级别:`PREFECT_LOGGING_LEVEL=WARNING`(或更高)
- 文档:Prefect 3 Settings、Logging 定制
时间分区(架构级,“删除变成丢分区”)
- 把 log 按 timestamp 做原生分区(日/周/月均可,按你的保留周期选择)。这样清理直接 `DROP PARTITION`,速度极快、几乎无写入放大。
- 典型迁移步骤(需评估窗口和回滚预案):
1) 建一个新的分区父表(结构同现有 `log`,按 timestamp RANGE 分区),创建默认分区和近期需要的分区;在分区上分别创建与现有相同的索引(父表上的索引需要在每个分区单独创建)。
2) 短暂停机或流量切换窗口内,重命名旧表 `log -> log_old`,把新的分区父表重命名为 `log`,确保约束/外键与原表一致。
3) 视情况把近一段时间的日志从 log_old 迁入对应分区(旧数据也可分批迁移或直接保留到过期时再按旧表删除)。
4) 之后的保留就变成:定时 ALTER TABLE log DETACH PARTITION log_2025_10; DROP TABLE log_2025_10;
- 如果不想自己维护分区管理,可以考虑 pg_partman 等工具。
脚本层面的改造建议
- 在你的例行作业最前面加一段“按时间删日志”的批处理循环,直到 0 行受影响为止。然后再执行你现有的 flow_run 相关清理(大概率 log 的两段 DELETE 可以去掉或极大缩小批量):
-- 0) 日志按时间分批清理(先执行)
WITH batch AS (
SELECT id
FROM log
WHERE timestamp < now() - interval '12 days'
ORDER BY timestamp
LIMIT 100000
)
DELETE FROM log l
USING batch
WHERE l.id = batch.id;
-- 循环上述块直到 0 行
- 保留你对 flow_run/task_run 的批处理和事务边界,保持每批 5k~20k,平衡锁与时长。
注意事项
- 确认 log.timestamp 的时区与 now() 一致(UTC 更安全),避免误删。
- 如果 log 上确实有外键到 flow/task(Prefect 3 的 ORM 字段允许为空),按时间直接删日志不会破坏参照完整性;但反向删除 flow/task 时,仍建议先按时间删日志,以免大批外键检查拖慢删除。
- 大表上强烈建议时间列用 BRIN 索引;对“以 ID JOIN 删除”的场景,确保 `log(flow_run_id)`、`log(task_run_id)` 有 BTree 索引,否则 JOIN 删除必扫全表。
相关文档
- Prefect 3 Settings(含 PREFECT_LOGGING_TO_API_ENABLED 等):Settings reference
- 日志模型与服务端 API:Logs REST API、Python models.logs
- 数据库维护指南(说明事件保留、日志表等):Database maintenance
- 日志定制:Logging customization
能否补充几条信息,我可以给出更精确的 DDL/迁移脚本建议:
- 你使用的 Postgres 版本?
- 是否可以接受把 log 表改为按时间分区?
- 是否可以降低日志保留期,或者降低/关闭把日志写入 Prefect API(UI 中就看不到历史日志了)?R G
11/20/2025, 2:59 AMMarvin
11/20/2025, 2:59 AMMarvin
11/20/2025, 3:01 AMFlowRunContext.get() 读取当前的 Flow 上下文,然后访问其中的 `flow_run.name`(也就是 `flow_run_name`)。
示例(同步 hook):
from prefect import flow, task
from prefect.context import FlowRunContext
def log_task_failure(task, task_run, state):
ctx = FlowRunContext.get()
if ctx and ctx.flow_run:
print(f"Flow run name: {ctx.flow_run.name}")
print(f"Flow run id: {ctx.flow_run.id}")
# 任务本身的信息也在参数里
print(f"Task run name: {task_run.name}")
print(f"Task run id: {task_run.id}")
@task(on_failure=[log_task_failure])
def my_task():
raise RuntimeError("boom")
@flow
def my_flow():
my_task()
if __name__ == "__main__":
my_flow()
补充说明:
- 如果极端情况下拿不到上下文(`FlowRunContext.get()` 返回 `None`,例如你在 hook 里开了自定义线程),也可以用 task_run.flow_run_id 通过 Orchestration 客户端读取 FlowRun:
from prefect.client.orchestration import get_client
async def log_task_failure_async(task, task_run, state):
async with get_client() as client:
flow_run = await client.read_flow_run(task_run.flow_run_id)
print(f"Flow run name: {flow_run.name}")
- hook 函数既可以是同步的,也可以是异步的;签名都是 `def/async def hook(task, task_run, state)`。
参考(源码/API):
- FlowRun 上下文与属性:prefect.context.FlowRunContext 源码R G
11/20/2025, 8:16 AMMarvin
11/20/2025, 8:16 AMMarvin
11/20/2025, 8:26 AMfrom prefect import flow, task
from prefect.futures import wait
from prefect.states import get_state_exception
from prefect.context import FlowRunContext
@task
def task_a(x): ...
@task
def task_b(x): ...
@task
def task_c(x): ...
@task
def task_d(tasks: dict): ...
@flow(flow_run_name="{logid}")
def my_flow(input_data, logid):
ctx = FlowRunContext.get()
flow_run_id = ctx.flow_run.id
# 提前用 dict 记录每个 future,方便失败时识别是哪个任务
tasks = {}
try:
a = task_a.submit(input_data)
tasks['a'] = a
b = task_b.with_options(retries=3).submit(a)
c = task_c.submit(a)
tasks['b'] = b
tasks['c'] = c
if a.result(raise_on_failure=True)["value"] > 30:
e = task_e.submit(a)
tasks['e'] = e
d = task_d.submit(tasks)
tasks['d'] = d
# 等全部完成,再统一检查状态(也可以按需早检查)
done, _ = wait(list(tasks.values()))
failed = {}
for name, fut in tasks.items():
if fut.state.is_failed() or fut.state.type.name in ("CANCELLED", "CRASHED"):
exc = get_state_exception(fut.state) # 不会抛出,直接取异常对象
failed[name] = {
"task_run_id": str(fut.task_run_id),
"state": fut.state.type.name,
"message": fut.state.message,
"exception": str(exc) if exc else None,
}
if failed:
return {
"flow_run_id": str(flow_run_id),
"errno": 1,
"errmsg": "some tasks failed",
"failed": failed,
}
return {
"flow_run_id": str(flow_run_id),
"errno": 0,
"errmsg": "",
"result": d.result(raise_on_failure=True),
}
except Exception as e:
# try/except 里也可以把已提交的任务做一次状态汇总
done, _ = wait(list(tasks.values())) if tasks else (set(), set())
failed = {}
for name, fut in tasks.items():
if fut.state.is_failed() or fut.state.type.name in ("CANCELLED", "CRASHED"):
exc = get_state_exception(fut.state)
failed[name] = {
"task_run_id": str(fut.task_run_id),
"state": fut.state.type.name,
"message": fut.state.message,
"exception": str(exc) if exc else None,
}
return {
"flow_run_id": str(flow_run_id),
"errno": 1,
"errmsg": str(e),
"failed": failed,
}
要点:
- 用 tasks = {'b': b, 'c': c, ...} 这样的 dict 记录每个提交的 future。
- 用 future.state.is_failed() 或 future.state.type 判断是否失败/取消/崩溃。
- 用 get_state_exception(future.state) 拿到异常对象,不会再次抛出。
- future.result(raise_on_failure=True) 会抛异常;如果只是想拿异常作为结果,传 `raise_on_failure=False`。
方式 B:让 flow 失败时,通过 on_failure hook 查询本次运行的失败任务
```
from prefect import flow
from prefect.client.orchestration import get_client
from prefect.client.schemas.filters import (
TaskRunFilter,
TaskRunFilterFlowRunId,
TaskRunFilterState,
TaskRunFilterStateType,
)
from prefect.client.schemas.objects import StateType
def list_failed_tasks_on_flow_failure(flow, flow_run, state):
# 同步客户端,便于在同步 hook 中使用
with get_client(sync_client=True) as client:
failed_task_runs = client.read_task_runs(
task_run_filter=TaskRunFilter(
flow_run_id=TaskRunFilterFlowRunId(any_=[flow_run.id]),
state=TaskRunFilterState(
type=TaskRunFilterStateType(any_=[StateType.FAILED, StateType.CRASHED, StateType.CANCELLED])
),
),
limit=500,
)
# 这里可以把失败任务名、错误信息等记录到日志或外部系统
for tr in failed_task_runs:
# tr.name 为任务运行名,tr.state.message 为错误信息
print(f"failed task: {tr.name} ({tr.id}) - {tr.state.type.value} - {tr.state.message}")Marvin
11/20/2025, 8:26 AMCANCELLED 或 `CRASHED`,可以一并统计。
相关源码参考:
- Futures/状态 API:prefect.futures 源码
- 客户端过滤器(TaskRunFilter 等):filters.py
- Orchestration 客户端:orchestration client
如果你愿意,我可以根据你现有的 task_a/b/c/d/e 具体返回结构,帮你把失败聚合的返回结构再细化一下。R G
11/20/2025, 8:56 AMMarvin
11/20/2025, 8:56 AMMarvin
11/20/2025, 8:56 AMget_state() 方法,直接用 future.state 读取状态。
- 建议先 wait() 一下,确保所有 future 都已完成,否则你读取到的可能是 PENDING/RUNNING。
- 如果你把“失败”定义得宽一点,通常会把 CRASHED/CANCELLED 也算进来(很多时候下游因上游失败而被取消)。
修正版(保留你用 locals() 扫描的写法)
from prefect.futures import PrefectFuture, PrefectFutureList, wait
from prefect.client.schemas.objects import StateType
failed_tasks_name = []
# 先收集所有 future,统一等它们完成
local_vars = locals()
all_futures = []
for var_val in local_vars.values():
if isinstance(var_val, PrefectFuture):
all_futures.append(var_val)
elif isinstance(var_val, (list, PrefectFutureList)):
all_futures.extend([f for f in var_val if isinstance(f, PrefectFuture)])
if all_futures:
wait(all_futures) # 等所有任务完成,避免状态还在 RUNNING/PENDING
# 再判断失败的变量名
for var_name, var_val in local_vars.items():
if isinstance(var_val, PrefectFuture):
st = var_val.state
if st.is_failed() or st.type in {StateType.CRASHED, StateType.CANCELLED}:
failed_tasks_name.append(var_name)
elif isinstance(var_val, (list, PrefectFutureList)):
if not var_val:
continue
# 只要列表里有一个失败/崩溃/取消,就把这个变量名记为失败集合
if any(
isinstance(fut, PrefectFuture) and
(fut.state.is_failed() or fut.state.type in {StateType.CRASHED, StateType.CANCELLED})
for fut in var_val
):
failed_tasks_name.append(var_name)
更推荐的做法(比扫描 locals 更稳)
- 在 flow 里维护一个明确的字典来追踪任务,键就是你想返回的“任务标识”,值是 future。
tasks = {}
tasks['b'] = task_b.with_options(retries=3).submit(a)
tasks['c'] = task_c.submit(a)
...
done, _ = wait(list(tasks.values()))
failed = [name for name, fut in tasks.items()
if fut.state.is_failed() or fut.state.type in {StateType.CRASHED, StateType.CANCELLED}]
补充
- 拿到异常对象而不抛出:`from prefect.states import get_state_exception; exc = get_state_exception(future.state)`
- 你现在收集的是“变量名”,不是任务运行名;如果想拿任务运行名,可以用客户端通过 future.task_run_id 读 TaskRun 再取 `name`。
- Prefect 3.x futures/状态源码参考:
- prefect.futures 源码
- states 源码
如果你愿意,我可以把这段逻辑直接嵌入到你现有的 my_flow 里,返回你期望的结构。R G
11/20/2025, 12:51 PMMarvin
11/20/2025, 12:51 PMMarvin
11/20/2025, 12:52 PM# 假设 fut 是 PrefectFuture
val = fut.result(raise_on_failure=False) # 失败时返回异常对象,不会抛错
if isinstance(val, Exception):
concise = str(val) # -> "blade_req_final_failed"
# 也可更保险:
if not concise and val.args:
concise = str(val.args[0])
方式 2:从状态取异常对象
from prefect.states import get_state_exception
st = fut.state
if st.is_failed() or st.is_crashed():
exc = get_state_exception(st) # 不会抛错,返回 Exception 或 None
concise = str(exc) if exc else (st.message or "")
在 hook 里同理:
from prefect.states import get_state_exception
def on_task_failure(task, task_run, state):
exc = get_state_exception(state)
concise = str(exc) if exc else (state.message or "")
print(concise) # "blade_req_final_failed"
补充说明:
- str(exc) 对 RuntimeError("blade_req_final_failed") 会得到纯消息字符串,不含堆栈。
- state.message 通常是 “RuntimeError: blade_req_final_failed”,如果你只想冒号后的部分,也可以做个简单切分:
concise = state.message.split(":", 1)[-1].strip() if state.message else ""
但这不如直接用 str(exc) 稳定。
- 避免用 `future.result(raise_on_failure=True)`,它会抛异常;用 raise_on_failure=False 拿到异常对象再取 str() 即可。
源码参考:
- get_state_exception:prefect.statesR G
11/21/2025, 2:52 AMMarvin
11/21/2025, 2:52 AMMarvin
11/21/2025, 2:53 AMtask_run_id 去读对应的 TaskRun,再取 `TaskRun.name`。下面是修正后的函数,也顺手修了你原来 list 分支里取 var_val.state 的小 bug,以及外层/内层循环的 break 逻辑。
说明:
- PrefectFuture 本身没有 task_name 属性;要么查询 Orchestration API 的 `TaskRun.name`(运行名),要么在提交任务时自己维护一个“任务定义名 → future”的映射。
- 我这里返回的是 Task Run 名称(通常就是你 Task 的可读名),以及精简的错误信息。
版本 1:直接查 TaskRun(最简单,零改动现有提交流程)
from prefect.futures import PrefectFuture, PrefectFutureList, wait
from prefect.client.orchestration import get_client
from prefect.client.schemas.objects import StateType
from prefect.states import get_state_exception
def get_failed_tasks(local_vars):
"""返回 (failed_task_run_name, concise_error_msg),如未发现失败则为空字符串"""
failed_task_name = ""
failed_task_msg = ""
# 可选:如果需要,先等待这些 futures 完成,避免读取到 RUNNING/PENDING
all_futs = []
for v in local_vars.values():
if isinstance(v, PrefectFuture):
all_futs.append(v)
elif isinstance(v, (list, PrefectFutureList)):
all_futs.extend([f for f in v if isinstance(f, PrefectFuture)])
if all_futs:
wait(all_futs)
def concise_err(st):
exc = get_state_exception(st)
return str(exc) if exc else (st.message or "")
with get_client(sync_client=True) as client:
def check_future(fut):
nonlocal failed_task_name, failed_task_msg
st = fut.state
if st.is_failed() or st.type in {StateType.CRASHED, StateType.CANCELLED}:
tr = client.read_task_run(fut.task_run_id)
failed_task_name = tr.name # 任务运行名
failed_task_msg = concise_err(st) # 精简错误
return True
return False
for var_val in local_vars.values():
if isinstance(var_val, PrefectFuture):
if check_future(var_val):
break
elif isinstance(var_val, (list, PrefectFutureList)):
for fut in var_val:
if isinstance(fut, PrefectFuture) and check_future(fut):
break
if failed_task_name:
break
return failed_task_name, failed_task_msg
版本 2:推荐做法(性能最佳):在提交时维护“任务定义名 → future”的映射
- 这样不需要额外 API 请求(不查 TaskRun),直接用你希望的任务名键。
from prefect.futures import wait
from prefect.client.schemas.objects import StateType
from prefect.states import get_state_exception
def collect_failed_from_mapping(tasks: dict[str, PrefectFuture]):
# tasks 形如:{"task_b": b_future, "task_c": c_future, ...}
if tasks:
wait(list(tasks.values()))
for name, fut in tasks.items():
st = fut.state
if st.is_failed() or st.type in {StateType.CRASHED, StateType.CANCELLED}:
msg = str(get_state_exception(st) or "") or (st.message or "")
return name, msg
return "", ""
示例用法:
tasks = {}
tasks["task_b"] = task_b.with_options(retries=3).submit(a)
tasks["task_c"] = task_c.submit(a)
# 如果有列表任务,建议用 "task_x[0]" 这类标记来存 key
failed_name, failed_msg = collect_failed_from_mapping(tasks)
小结与注意
- 你原代码里的 failed_task_msg = get_err_msg(var_val.state) 在 list 分支应改为 `get_err_msg(fut.state)`。
- 是否把 CANCELLED/CRASHED 也算失败,按你的业务需要调整集合 `{StateType.CRASHED, StateType.CANCELLED}`。
- 如果你确实需要“任务定义名”(装饰器里的 name),最佳方式是在提交时就把 task_obj.name 当作字典 key 存起来;从 future 反查定义名在 Prefect 3 里没有直接属性。
相关源码参考:
- Futures 与状态:prefect.futures
- 状态与异常提取:prefect.states
- Orchestration 客户端:orchestration clientR G
11/21/2025, 4:55 AMMarvin
11/21/2025, 4:55 AMMarvin
11/21/2025, 4:56 AMget_client 指向的 API 与当前 flow run 的 API 不同
- 极少数情况下,ID 不是 task run(比如 subflow 的 flow run future)
为了稳妥拿到“任务名”,我给你两套可用方案:
方案 A(推荐):提交时维护“任务名 → future”的映射,不依赖 API
- 优点:不需要调用 API,离线/本地都稳定
- 你可以用 task_obj.name 作为键,这就是任务定义名(接近你想要的“task名”)
from prefect.futures import PrefectFuture, PrefectFutureList, wait
from prefect.client.schemas.objects import StateType
from prefect.states import get_state_exception
def get_failed_tasks_from_map(tasks_map: dict[str, PrefectFuture] | None,
lists_map: dict[str, list[PrefectFuture]] | None = None):
"""
tasks_map 例子: {"task_b": b_future, "task_c": c_future}
lists_map 例子: {"task_x": [x_fut0, x_fut1, ...]}
返回 (failed_task_name, concise_error_msg)
"""
tasks_map = tasks_map or {}
lists_map = lists_map or {}
# 等待所有提交完成(避免读到 RUNNING/PENDING)
all_futs = list(tasks_map.values()) + [f for lst in lists_map.values() for f in lst]
if all_futs:
wait(all_futs)
def concise_err(st):
exc = get_state_exception(st)
return str(exc) if exc else (st.message or "")
# 先查单个任务
for name, fut in tasks_map.items():
st = fut.state
if st.is_failed() or st.type in {StateType.CRASHED, StateType.CANCELLED}:
return name, concise_err(st)
# 再查列表任务(任一元素失败就返回该列表名)
for name, futs in lists_map.items():
for fut in futs:
st = fut.state
if st.is_failed() or st.type in {StateType.CRASHED, StateType.CANCELLED}:
return name, concise_err(st)
return "", ""
用法示例:
tasks_map = {}
tasks_map[task_b.name] = task_b.with_options(retries=3).submit(a)
tasks_map[task_c.name] = task_c.submit(a)
lists_map = {}
# 如果有批量任务,给它一个集合名
# lists_map["task_x"] = [task_x.submit(i) for i in items]
failed_name, failed_msg = get_failed_tasks_from_map(tasks_map, lists_map)
方案 B:继续用 API 查询 TaskRun.name,但要容错和兜底
- 确认你的 flow 运行时已连接到 Prefect API(设置好 `PREFECT_API_URL`/profile)
- 在函数里对 ObjectNotFound 做捕获,失败时退回到方案 A 或变量名
这是你原函数的修正版(修了 list 分支中对 var_val.state 的 bug,并增加了 API 查询 + 兜底):
```
from prefect.futures import PrefectFuture, PrefectFutureList, wait
from prefect.client.orchestration import get_client
from prefect.client.schemas.objects import StateType
from prefect.states import get_state_exception
from prefect.exceptions import ObjectNotFound
def get_failed_tasks(local_vars, tasks_map: dict[str, PrefectFuture] | None = None):
"""返回 (failed_task_run_name 或 任务键名, concise_error_msg)"""
# 先等待 locals 里所有 futures 完成
all_futs = []
for v in local_vars.values():
if isinstance(v, PrefectFuture):
all_futs.append(v)
elif isinstance(v, (list, PrefectFutureList)):
all_futs.extend([f for f in v if isinstance(f, PrefectFuture)])
if all_futs:
wait(all_futs)
def concise_err(st):
exc = get_state_exception(st)
return str(exc) if exc else (st.message or "")
def name_from_map(fut):
if not tasks_map:
return None
for k, v in tasks_map.items():
if v is fut:
return k
return None
def name_from_api(fut):
try:
with get_client(sync_client=True) as client:
tr = client.read_task_run(fut.task_run_id)
return tr.name
except ObjectNotFound:
return None
except Exception:
return None
# PrefectFuture
for var_name, var_val in local_vars.items():
if isinstance(var_val, PrefectFuture):
st = var_val.state
if st.is_failed() or st.type in {StateType.CRASHED, StateType.CANCELLED}:
msg = concise_err(st)
# 优先 API,其次任务映射,最后退回变量名
name = name_from_api(var_val) or name_from_map(var_val) or var_name
return name, msgMarvin
11/21/2025, 4:56 AMmy_flow 里,帮你把返回结构也统一好。R G
11/21/2025, 7:14 AMMarvin
11/21/2025, 7:14 AMMarvin
11/21/2025, 7:15 AMsetdefault 赋给 name 和 task_run_name 即可(用户手动传入时仍能覆盖)。
给你一版修正且支持 @task 和 @task(...) 的实现(顺便修了你变量名的小问题):
from prefect import task as prefect_task
def task(__fn=None, **kwargs):
# 统一默认
kwargs.setdefault("cache_result_in_memory", False)
kwargs.setdefault("persist_result", False)
kwargs.setdefault("cache_policy", NO_CACHE)
kwargs.setdefault("on_failure", [log_task_failure])
def _decorate(fn):
# 默认让 Task 定义名 与 Task Run 名称 都等于函数名
kwargs.setdefault("name", fn.__name__)
kwargs.setdefault("task_run_name", fn.__name__)
return prefect_task(**kwargs)(fn)
# 支持 @task 和 @task(...)
return _decorate if __fn is None else _decorate(__fn)
说明与注意:
- name 是任务定义名;`task_run_name` 是每次运行时的运行名。上面写法让两者默认都等于函数名。
- 若你用 with_options(task_run_name=...) 或在装饰器入参显式传入,会覆盖默认。
- 如果你后续想让运行名更有区分度(例如批量提交多次时加上参数),也可以把 task_run_name 换成可调用:
kwargs.setdefault(
"task_run_name",
lambda ctx, params: fn.__name__ # 这里可以根据 params 组装更丰富的名字
)R G
11/21/2025, 7:20 AMMarvin
11/21/2025, 7:20 AMMarvin
11/21/2025, 7:22 AMTaskRun 里没有“函数名”这个字段。能用的是 `task_key`,它的格式是“函数名-代码哈希”。因此你可以从 task_run.task_key 里解析出函数名部分。
实用做法
- 直接解析 `task_key`:
def extract_func_name_from_task_key(task_key: str) -> str:
# task_key 形如 "my_task-a1b2c3d4" 或 "package.module.MyClass.my_task-a1b2c3d4"
# 先去掉模块/类限定,再去掉哈希
return task_key.split(".")[-1].split("-")[0]
- 放到你的查询代码里用:
from prefect.client.schemas.filters import TaskRunFilter, TaskRunFilterFlowRunId
task_runs = await client.read_task_runs(
task_run_filter=TaskRunFilter(
flow_run_id=TaskRunFilterFlowRunId(any_=[flow_run.id])
)
)
func_names = [extract_func_name_from_task_key(tr.task_key) for tr in task_runs]
补充说明
- task_run.name 是“任务运行名”(未显式设置时是随机 slug),不是函数名。
- Prefect 3.x 并不在 API 层保存 Task 定义对象,只有 TaskRun 记录;因此没有“直接返回函数名”的字段。
- 如果你希望完全控制“任务运行名 = 函数名”,也可以在装饰器里把 task_run_name 默认设成 `fn.__name__`(你之前那段装饰器方案即可),那样在 API 侧就能通过 task_run.name 看到函数名。但这属于“运行名”,不是“函数定义名”。
如果你愿意,我可以把上面的 extract_func_name_from_task_key 集成到你现有的 failed-tasks 统计逻辑里,返回更友好的“函数名 + 精简错误信息”。R G
11/21/2025, 8:10 AMMarvin
11/21/2025, 8:10 AMMarvin
11/21/2025, 8:10 AMstate 只有在任务进入终态(COMPLETED/FAILED/CRASHED/CANCELLED)后才会反映最终结果;未完成时就是 PENDING/RUNNING。
解决思路有三种,任选其一即可:
方案 A:先等待所有 futures 完成,再检查状态(推荐)
from prefect.futures import PrefectFuture, PrefectFutureList, wait
from prefect.client.schemas.objects import StateType
def get_failed_tasks(local_vars, timeout=None):
print("True Order:", list(local_vars.keys()))
# 先收集并等待
all_futs = []
for v in local_vars.values():
if isinstance(v, PrefectFuture):
all_futs.append(v)
elif isinstance(v, PrefectFutureList):
all_futs.extend([f for f in v if isinstance(f, PrefectFuture)])
done, not_done = wait(all_futs, timeout=timeout) if all_futs else (set(), set())
def is_failed(st):
return st.is_failed() or st.type in {StateType.CRASHED, StateType.CANCELLED}
failed_var = ""
failed_msg = ""
for var_name, var_val in local_vars.items():
if isinstance(var_val, PrefectFuture):
st = var_val.state
if is_failed(st):
failed_var = var_name
failed_msg = get_err_msg(st)
break
elif isinstance(var_val, PrefectFutureList):
if not var_val:
continue
for fut in var_val:
if not isinstance(fut, PrefectFuture):
continue
st = fut.state
if is_failed(st):
failed_var = var_name
failed_msg = get_err_msg(st)
break
if failed_var:
break
# 可选:如果有未完成任务,按需视为超时/失败
if not failed_var and not_done:
failed_var = "pending_tasks"
failed_msg = f"{len(not_done)} task(s) not finished"
return failed_var, failed_msg
方案 B:在读取状态前,按需阻塞到该 future 完成
- 对单个元素:
val = fut.result(raise_on_failure=False) # 等到完成,不会抛异常;失败时返回异常对象
st = fut.state # 此时为终态
- 对列表:
for fut in var_val:
fut.result(raise_on_failure=False)
if fut.state.is_failed():
...
方案 C:用 as_completed 按完成顺序扫描,遇到第一个失败就返回
from prefect.futures import as_completed
from prefect.client.schemas.objects import StateType
def get_first_failed(futs):
for fut in as_completed(futs, timeout=None):
st = fut.state
if st.is_failed() or st.type in {StateType.CRASHED, StateType.CANCELLED}:
return fut
return None
补充建议
- 扫描 locals 容易误伤其他变量,建议提交时维护一个专门的映射,比如:
tasks = {}
tasks["task_b"] = task_b.submit(...)
tasks["task_c"] = task_c.submit(...)
done, _ = wait(list(tasks.values()))
failed = [name for name, fut in tasks.items() if fut.state.is_failed()]
- 获取精简错误可继续用你现有的 `get_err_msg(st)`,或用 from prefect.states import get_state_exception 再 `str(exc)`。
参考源码:
- futures.wait / as_completed:prefect.futures
- 状态与异常提取:prefect.statesR G
11/25/2025, 11:00 AMMarvin
11/25/2025, 11:00 AMMarvin
11/25/2025, 11:04 AMprefect.futures.wait() 使用事件驱动的等待机制:进程内会有一个共享的异步事件循环/订阅器,而不是为每个 future 新建一个线程。对成千上万个 futures 调用一次 wait() 也不会产生成比例的新线程。
- 不会造成持续的 PostgreSQL 压力
- wait() 对远程运行的任务会做一次“初次读取”检查(每个 future 至少一次 API 读取任务状态),之后通过事件流(websocket)等待状态变化,不做高频轮询。
- 访问 future.state 时,如果状态未缓存且未终结,会触发一次按需的 API 读取;状态一旦终结,会在本地缓存,之后反复访问不会再次打 API。
- 真正会放大压力的反例
- 频繁循环访问 .state 并 sleep 的做法会造成高频 API 读取(类似轮询),应避免:
# 反例:会造成高频 API 调用
while not fut.state.is_final():
time.sleep(0.1)
你这段代码的影响评估
- 先收集 `PrefectFuture`/`PrefectFutureList`,再统一 `wait(all_futs, timeout=60)`,然后遍历检查 `state.is_failed()`:
- 这是推荐做法之一:一次性 wait 可以复用同一个事件订阅,避免 per-future 阻塞/轮询。
- timeout=60 后如果仍有未完成的 futures,接下来对这些未完成的 futures 访问 .state 时,可能会触发按需的 API 读取(一次/个),但不会形成高频轮询。
- 遍历 locals 去找 futures 本身没问题,但更建议在创建任务时自己维护一个 list,避免无关变量混入。
更高效/更稳妥的写法示例
- 一次性等待 + 批量取状态(避免不必要的结果读取):
from prefect.futures import wait, resolve_futures_to_states
def get_failed_tasks(local_vars):
failed = []
# 收集 PrefectFuture
futs = []
for v in local_vars.values():
if isinstance(v, PrefectFuture):
futs.append(v)
elif isinstance(v, PrefectFutureList):
futs.extend([f for f in v if isinstance(f, PrefectFuture)])
# 一次性等待(事件驱动,无线程暴增)
if futs:
wait(futs, timeout=60)
# 批量解析状态(避免逐个触发多次 API)
states = resolve_futures_to_states(futs) if futs else []
for fut, st in zip(futs, states):
if st.is_failed():
# 按需获取错误信息:如果你有自定义的 get_err_msg(st) 可以放这里
failed.append((fut, st))
return failed
- 如果你需要逐个尽快发现失败(不等全部完成),可以用 as_completed 逐步消费,但同样不会引入线程暴涨。
操作建议
- 保留“一次性 wait 再检查”的模式,避免 while 轮询 `.state`。
- 在单个 flow 里 futures 数量极大(例如几万级)时,初次状态读取会带来一次性的 API 峰值,相应地也会间接打到 Postgres。可以通过:
- 限制任务并发(TaskRunner 并发/映射规模控制)
- 分批 `wait`(分段提交/分段等待)
来平滑这一波峰值。
- 只在需要时调用 `state.result()`(它可能会触发结果存储的读取),平时用 `state.is_failed()`/`state.message` 足以定位失败。
如果你方便的话,补充两点信息我可以给更具体的优化建议:
- 单个 flow 里大概有多少个任务/futures?
- 这些任务是本地执行还是通过远程 worker 执行(比如 Docker/K8s/AWS 等)?你目前是否已经观察到 Postgres 的压力峰值,还是担心潜在风险?