<@ULVA73B9P> When using Prefect 3 client-side mode...
# ask-marvin
r
@Marvin When using Prefect 3 client-side mode to execute a flow, can a single process only run one flow at a time?
m
thought for 786.9 seconds
Short answer: Not strictly. In Prefect 3 client-side runs, calling a flow (e.g.,
my_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
Copy code
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)
Copy code
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
@Marvin I have developed a service that receives requests one by one and then calls Prefect 3 flows to process data locally before returning the results. Since the main program is already running using
asyncio.run
, is it necessary to use
<http://asyncio.to|asyncio.to>_thread
when invoking the flows?
m
thought for 234.9 seconds
Short answer: - If your flow is async (
@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
Copy code
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
Copy code
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)
Copy code
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
@Marvin I have developed a service that receives requests one by one and calls Prefect 3 flows to process data locally in client-side mode before returning the results. My flow is a long-running process that takes an average of 30 minutes to complete. I am facing an issue where memory usage continues to rise. Are there any solutions to this problem?
m
thought for 287.5 seconds
Yes—there are several practical ways to stop memory from creeping up during long client-side runs. Why this happens - By default, Prefect caches task/flow results in-memory for fast reuse. Over a 30-minute flow, especially if you create many tasks or large results, memory can trend upward. - Large objects returned from tasks/flows stay referenced until the flow ends. - Third-party libraries (e.g., pandas, NumPy) can fragment memory; even after del/GC, RSS may not return to the OS. - Heavy logging can buffer in memory. Recommended mitigations 1) Turn off in-memory result caching - At the flow and task level:
Copy code
from 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:
Copy code
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:
Copy code
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:
Copy code
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
1. @Marvin The flow contains over 20 task operations. 2. The intermediate results of the tasks within the flow are only used internally, and only the result of the final task needs to be returned via the flow return for the caller to obtain the final result. 3. All tasks are I/O-bound scenarios.
m
thought for 116.3 seconds
Thanks for the details — given: - ~20 tasks per run - Only the final result matters to the caller - All tasks are I/O-bound The most effective fixes to stop memory from creeping up are: 1) Disable in-memory result caching and persistence - Since you don’t need intermediate results after they’re used, turn off Prefect’s in-memory cache and persistence at both the flow and task level.
Copy code
from 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.
Copy code
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.
Copy code
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:
Copy code
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
@Marvin How to use
run_flow_in_subprocess
specifically and how to obtain the flow's running results.
m
thought for 66.1 seconds
Great question.
run_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)
Copy code
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)
Copy code
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)
Copy code
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
@Marvin Is it possible to avoid using
run_flow_in_subprocess
and directly use
multiprocessing
instead?
m
thought for 53.1 seconds
Yes, you can. You don’t have to use run_flow_in_subprocess; you can start a separate process yourself with multiprocessing and call the flow inside it. Just be mindful of a few caveats. Recommended pattern (portable and safe) - Use the “spawn” start method (important on Unix to avoid fork-related issues). - Put the flow call in a top-level function. - Return the final result via a multiprocessing.Queue or write it to a file (JSON/parquet/etc.). - Guard your entrypoint with if name == "__main__": to support Windows/macOS. Example using a Queue
Copy code
# 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)
Copy code
# 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)
Copy code
# 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)
Copy code
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
@Marvin Provide the usage method and examples for
run_sync_in_worker_thread
.
m
thought for 204.8 seconds
Here’s how to use
run_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:
Copy code
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
Copy code
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
Copy code
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
Copy code
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
Copy code
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.
- `run_async_in_new_loop`: Spins up a brand new event loop to run async code. Use this when you explicitly need an isolated loop; otherwise, prefer
run_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 documentation
r
@Marvin x1 = task1.map(x0); x2 = task2.submit(x1) Does the
map
syntax here guarantee that all tasks in task1 are completed successfully before executing task2?
m
thought for 405.0 seconds
Short answer: Yes. In Prefect 3.x, passing the result of
task1.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
Copy code
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
Copy code
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
Copy code
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:
Copy code
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
@Marvin from prefect import task as prefect_task from prefect.cache_policies import NO_CACHE def task(*args, **kwargs): """task""" # 统一设定部分参数默认值 kwargs.setdefault("cache_result_in_memory", False) kwargs.setdefault("persist_result", False) kwargs.setdefault("cache_policy", NO_CACHE) return prefect_task(*args, **kwargs) 我继承prefecr3 task装饰器自定义了实现,我现在想打印输入输出到控制台,怎么写代码
m
thought for 256.9 seconds
你可以有两种常用方式在 Prefect 3.x 里把任务的“入参与出参”打印到控制台/日志里: 1) 在任务函数里直接 print,并在装饰器上开启 log_prints - 最简单,也最直观。但要确保 `log_prints=True`,这样 print 会被 Prefect 捕获并出现在 Prefect 的日志里(本地运行也会直接到 stdout)。 示例(结合你的自定义装饰器,默认开启 log_prints):
Copy code
from 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 集成到你自定义装饰器里的完整例子(自动打印输入与输出):
Copy code
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
@Marvin 以下是我的flow其中一部分代码,帮我看看写的是否合理,不合理的话怎么修改。代码: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) if len(aigc_video_res_list) == 0: video_effect_res = blades_req.submit(video_effect_predata) keep_ori_audio = False else: video_effect_res = video_effect_predata keep_ori_audio = True merge_predata = preproc_merge.submit(video_effect_res, voice_effect_res, keep_ori_audio) merge_res = blades_req.submit(merge_predata)
m
thought for 314.4 seconds
整体思路是合理的:用
.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()`,会阻塞并降低并行度。你当前代码没有这么做,挺好。 建议的重写(等价但更紧凑):
Copy code
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:
Copy code
keep_ori_audio = bool(aigc_video_res_list.result())
或者写一个小 task:
Copy code
@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
之前做一次规整:
Copy code
@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
@Marvin 我自己部署的prefect server,client提交flow执行时返回500 Internal Server Error,请问怎么排查是什么问题?
m
thought for 515.8 seconds
先帮你快速定位思路:500 一般是服务端在处理请求时抛出了异常(而不是请求参数校验问题,后者通常是 422)。排查要点是“确认服务是否健康 → 看服务端日志 → 核对版本/数据库/服务组件 → 复现实验”。 下面给你一套 Prefect 3.x 自托管 Server 的检查清单,按顺序执行即可。执行命令前,确保你的客户端和 Server 所在环境都能运行
prefect
CLI。 - 文档参考(建议随查随对照) - Server 概念 - Server CLI - 健康检查接口 - 日志与调试 - 自托管参考架构 一、快速健康检查(API 与就绪) - 确认你客户端使用的
PREFECT_API_URL
指向的是 Server 的
/api
路径,例如
http://<host>:4200/api
- 健康/就绪探针
Copy code
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 的栈:
Copy code
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 与服务进程分离运行,也可单独跑服务:
Copy code
prefect server services start
并关注该进程的日志(调度、触发器、事件持久化等后台服务出错也会连带 500) 三、版本一致性检查(常见 500 根因) - 新客户端对旧 Server 发请求时,若 API 发生不兼容变更,可能触发 500
Copy code
# 分别在客户端和服务端环境执行
  prefect version
若客户端版本较新,建议先升级服务端到同一大版本的最新补丁(或对齐版本)再试。 四、数据库检查与迁移 - 查看当前配置,确认 Server 连接到了你预期的数据库:
Copy code
prefect config view --show-defaults | grep -i -E "DATABASE|API_URL|SERVER"
- 执行数据库迁移(很多 500 是因为表结构未升级/不完整)
Copy code
prefect server database upgrade -y
如遇大型库迁移超时,可提高超时:
Copy code
PREFECT_SERVER_DATABASE_TIMEOUT=30 prefect server database upgrade -y
- 使用 PostgreSQL 时,确保安装了
pg_trgm
扩展(Prefect 用于搜索/索引;缺失可能导致查询相关的 500)
Copy code
CREATE EXTENSION IF NOT EXISTS pg_trgm;
如果是托管 PG,请在 DB 管理端确认是否支持该扩展。 - 若是 SQLite,确认路径可写、未被多个进程/容器并发写(并发写会报错) 五、后台服务状态(调度/触发器等) - 列出并检查服务状态:
Copy code
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
@Marvin 我使用docker部署的prefect server,进入docker后怎么查看prefect server的输出日志
m
thought for 37.3 seconds
最简单可靠的方式是从宿主机看容器的 stdout/stderr;Prefect Server 默认不会把日志写成文件。 常用方法 - 宿主机直接看日志
Copy code
# 先找容器名
  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 部署
Copy code
# 查看所有服务日志
  docker compose logs -f

  # 只看某个服务(例如 api / services)
  docker compose ps
  docker compose logs -f <service_name>
在容器内查看(不推荐,但可以) - 进入容器后,直接跟主进程(通常 PID 1)的 stdout/stderr
Copy code
# 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 或启动命令后重启容器)
Copy code
# 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
@Marvin flow内部能否通过trycache获取到失败的task是哪个
m
thought for 128.9 seconds
可以,但要注意捕获的位置。在 Prefect 3.x 里,`@task` 在 flow 中调用时会返回一个 PrefectFuture;真正的异常是在你取结果时才会抛出。因此需要对 `future.result()`(或先转 state 再取异常)做 try/except,才能知道哪个 task 失败。 常用几种写法 1) 直接对 result() 做 try/except,并保留上下文信息(比如参数或索引)
Copy code
from 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 判断
Copy code
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 对象
Copy code
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
@Marvin 针对上一个问题我补充一下,flow可能执行了大量task,而我没法逐个检查全部的task返回的future对象。比如我通过在task里return Failed(message=xxx)是否可以让flow直接捕获失败的task名字
m
thought for 100.9 seconds
结论先说: - 不建议在 task 里 `return Failed(message=...)`。在 Prefect 3.x 里,这只会把
Failed
这个对象当作“正常返回值”,不会把 task 状态标记为失败,更不能让 flow “自动知道哪个 task 失败了”。如果要让 task 真的失败,请在 task 内部直接
raise
异常。 - 想在不逐个
future.result()
的情况下“批量”拿到失败任务,可以在提交时统一收集 future,最后一次性解析所有 state 来筛选失败项;或者让任务自己在失败时记录名字。 可选做法一:统一收集 + 批量解析状态(推荐)
Copy code
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(...)
能拿到任务名等元数据。 可选做法二:任务内自报姓名并失败(无需集中扫描)
Copy code
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”)
Copy code
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
@Marvin 我搭建了自运维的prefect server,我发现pgsql的存储空间占用上涨非常快,快到存储极限了,怎么解决?我的pgsql存储空间总共3T,每天处理flow有50w。
m
thought for 648.3 seconds
看起来你是自托管的 Prefect Server(假设使用 Prefect 3.x),且每天大约有 50 万个 flow run。PostgreSQL 增长过快通常有这几个主要来源:events、logs,其次是 flow/task runs、states、artifacts。下面给你一套“立刻止血 + 有序清理 + 长期优化”的方案。 先给出结论(按优先级,从影响最大的开始) 1) 大幅减少新数据的写入量(立刻见效) - 关闭或减少日志上报到 API - 将运行环境中的 `PREFECT_LOGGING_TO_API_ENABLED=false`(建议优先做) - 或把日志级别提高到 `WARNING`/`ERROR`,减少 INFO 级日志:`PREFECT_LOGGING_LEVEL=WARNING` - 可调小单条日志和批量上报阈值:`PREFECT_LOGGING_TO_API_MAX_LOG_SIZE`、`PREFECT_LOGGING_TO_API_BATCH_SIZE`、`PREFECT_LOGGING_TO_API_BATCH_INTERVAL` - 文档:Logging customization - 设置参考:Settings reference - 缩短 events 保留期 - 将 `PREFECT_SERVER_EVENTS_RETENTION_PERIOD`(或别名 `PREFECT_EVENTS_RETENTION_PERIOD`)设为更短,例如
P1D
或 `P3D`(ISO8601 时长,默认是 `P7D`) - 这对数据库增长影响巨大,因为 events 表是增长最快的 - 文档:Settings reference 2) 清理已有的历史数据(释放磁盘空间) - 官方没有内置“一键批量清理”CLI,需要用数据库维护或脚本化清理 - 推荐顺序:events → logs → flow/task runs → artifacts - 维护指南:Database maintenance A. 先找出最大表(诊断) 你可以先在 PostgreSQL 中确认体量最大的表,再决定清理策略:
Copy code
-- 看前 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):
Copy code
-- 预览将要删除的事件范围(举例:删除 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 清理:
Copy code
-- 预览
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 验证,列名以实际为准):
Copy code
-- 预览 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 ~ P3D
- flow/task runs:根据审计合规与可观测性需求,常见是 7~30 天;更长建议归档到外部存储而非数据库 - 日志外部化 - 将运行日志导向 stdout 并由外部日志系统(如 Loki/ELK/Splunk)采集,Prefect DB 不再存日志 - 为 events/logs 做时间分区(高级) - 如果你必须长保留,考虑对大表做基于时间的分区,配合定期 DROP PARTITION;Prefect Cloud 团队在生产上采用过分区方案,参考: - 在生产 Postgres 上分区而不停机 - 建立“每日清理”维护作业 - 用 Cron 或一个 Prefect 维护 flow,每天执行清理 SQL + VACUUM,保持稳定的占用和统计信息 4) Docker Compose/环境变量示例 如果你使用 docker-compose 启动 server/worker,可以在相关服务中添加:
Copy code
environment:
  # 强烈建议:减少或关闭日志上报
  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
@Marvin prefect config set PREFECT_EVENTS_RETENTION_PERIOD="2d" 这样设置能有效吗
m
thought for 55.9 seconds
短答:你这条命令里的值格式不对;要用 ISO8601 时长格式。把
2d
改成
P2D
才会生效。此外,最好用服务端键名,并确保是在运行 Prefect Server 的同一环境里设置并重启服务。 建议用法 - 正确设置(建议使用服务端键名):
Copy code
prefect config set PREFECT_SERVER_EVENTS_RETENTION_PERIOD="P2D"
- 你用的键名
PREFECT_EVENTS_RETENTION_PERIOD
在 3.x 中有别名支持,一般也能识别,但推荐显式使用 `PREFECT_SERVER_EVENTS_RETENTION_PERIOD`。 - 设置后执行校验,确认没有解析错误:
Copy code
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
@Marvin 我发现pgsql占用存储最多的表分别是events(219G)、task_run_state(157G)、event_resources(112G)、task_run(105G)、log(69G)、flow_run(1G)、flow_run_state(0.5G),怎么清理合适
m
thought for 105.0 seconds
谢谢你给出具体表大小,这很有帮助。结合这个分布(events 219G、task_run_state 157G、event_resources 112G、task_run 105G、log 69G、flow_run 1G、flow_run_state 0.5G),建议按下面顺序执行“限流-清理-维护”三步走。以下示例默认你在 Prefect 3.x,自托管 PostgreSQL,且允许直接执行 SQL。 一、先“止血”:减少新增写入 - 缩短 events 保留期:把服务端环境变量设为较短窗口(例如 1~3 天) -
PREFECT_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 清理)
Copy code
-- 哪些表引用了某个表(例如 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 天前)
Copy code
-- 预览
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 清理(根据你的外键结构调整联结):
Copy code
-- 示例:如果 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 - 如果已经决定关闭或缩短日志保留,这里一次性清理历史日志
Copy code
-- 预览
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 天为准;若有未结束的运行请避开):
Copy code
-- 先找出要删除的旧 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) 执行维护(释放空间) - 对大表执行:
Copy code
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 天
- flow/task runs:7~30 天(依合规需求),超长保留建议外部归档 - 日志外部化:让 Worker/Flow 的日志走 stdout,由 Loki/ELK/Splunk 收集,Prefect DB 不存日志 - 定期自动清理:用 Cron 或 Prefect 自身的维护 flow,每天跑一次“时间窗口 + 分批删除 + VACUUM/pg_repack” - 超长期或高留存:考虑对 events/logs 做时间分区+按分区 DROP - 参考:在生产 Postgres 上分区而不停机 需要你确认的点(好给你更精确的 SQL/脚本) - 你能接受的各类数据保留期(events、logs、flow/task runs) - 是否可以直接在 PostgreSQL 执行 SQL(以及是否能安装 pg_repack) - Prefect Server 是否运行在容器里,方便我们给出环境变量和重启步骤 - 对 task_run_state 是否有“只保留最后状态”的保留需求,还是可以随旧 flow 整批删除 如果你提供上述信息,我可以按你的保留期为你生成可直接执行的“分批清理 SQL 脚本”和“每日维护作业”模板。文档再附一次: - 设置参考:Settings reference - 日志定制:Logging customization - 数据库维护:Database maintenance
r
@Marvin 对于以上问题我想使用python prefect3 client清理存储资源,怎么实现?
m
thought for 356.3 seconds
可以用 Prefect 3 的 Python Client 做“按保留期批量清理”的脚本,优先清理 flow runs(及其关联对象,如 task runs / states)、artifacts。注意两点: - events / event_resources 没有对外的删除 API,需依赖服务端保留期自动清理或直接 SQL;用客户端无法删除它们。 - logs 目前也没有公开的“批量删除”API;删除 flow run 是否会级联清日志不保证,因此如果必须清理大量历史日志,建议用数据库层面清理。 下面给你一份可直接运行的 Sync 客户端清理脚本(推荐用于 cron/运维作业)。它会: - 找出早于保留期且已结束的 flow runs(通过 start_time 和状态筛选) - 先删这些 run 下的 artifacts(如果有) - 再删 flow run(通常会级联 task runs 和 states) 在生产前,先用 DRY_RUN 模式验证范围。 示例脚本(Sync)
Copy code
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://<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()
要点和注意
- 你给出的最大表里,events 和 event_resources 体量很大:客户端无法删它们。请务必同时设置服务端事件保留期缩短数据增长: - 在运行 Prefect Server 的环境/容器设置:`PREFECT_SERVER_EVENTS_RETENTION_PERIOD="P2D"`(或更短,ISO8601 时长) - 文档:Settings reference - logs 目前没有公开的“按过滤器删除”客户端方法;删除 flow run 不保证会清理其历史日志。如需清理大量历史日志,建议还是用数据库层面按 timestamp 分批删除,或干脆关闭日志上报到 API,避免继续涨: -
PREFECT_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
@Marvin 你给的命令不对:prefect=> DELETE FROM log WHERE "timestamp" < now() - interval '30 days' ORDER BY "timestamp" LIMIT 100000; ERROR: syntax error at or near "ORDER" LINE 3: ORDER BY "timestamp"
m
thought for 34.3 seconds
你说得对,PostgreSQL 不支持在 DELETE 语句里直接用 ORDER BY / LIMIT。需要用 CTE 或子查询来“先选出一批要删的主键”,再执行 DELETE。 修正示例(按主键 id 批删,推荐)
Copy code
-- 建议先确保有时间列索引,加速筛选
-- 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(单语句内唯一):
Copy code
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 万行,直到删不动为止)
Copy code
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 时间):
Copy code
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
@Marvin 分别都保留15天
m
thought for 76.8 seconds
明白,保留期统一为 15 天。基于你当前表体积分布,推荐组合方案: - 用 Python Prefect 3 Client 清理 flow runs(连带删掉其 artifacts,通常会级联清理 task runs 和 states) - 用 SQL 批量清理 events、event_resources、logs(这三类没有对外的批量删除 API) 同时把服务端保留策略/日志写入做“止血”: - 设置事件保留:`PREFECT_SERVER_EVENTS_RETENTION_PERIOD="P15D"` - 关闭或减少日志上报到 API:`PREFECT_LOGGING_TO_API_ENABLED=false` 或至少
PREFECT_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 验证
Copy code
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 再删
- 大表删除前建议先给时间列建索引(CONCURRENTLY,减少锁) - 批量执行,重复运行直到删不动为止 - 删除大量数据后记得 VACUUM/ANALYZE 或用 pg_repack 1) 预先索引(如未建立)
Copy code
-- 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 天前)
Copy code
-- 重复执行,直到 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 的联结。
Copy code
-- 删掉不再被任何 event 引用的资源行
DELETE FROM event_resources er
WHERE NOT EXISTS (
  SELECT 1 FROM event e WHERE e.resource_id = er.id
);
4) logs(删除 15 天前)
Copy code
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 为依据,只处理已结束的运行;按照“依赖→主记录”的顺序删除,避免外键约束失败。
Copy code
-- 先挑出旧的 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) 维护(释放空间)
Copy code
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 的环境设为
Copy code
PREFECT_SERVER_EVENTS_RETENTION_PERIOD="P15D"
- 日志上报:尽量关闭或降级
Copy code
PREFECT_LOGGING_TO_API_ENABLED=false
# 或
PREFECT_LOGGING_LEVEL=WARNING
参考: - 设置参考:Settings reference - 日志定制:Logging customization 需要我帮你把上述 SQL 打包成可重复执行的脚本(例如 psql 可执行 + 循环批删),或者把 Python Client 版换成异步并增加并行度的话,告诉我你偏好的运行方式(裸机/容器、是否可安装 pg_repack、可接受的批次大小/限流)。我可以按你的环境生成一份一键脚本。
r
@Marvin 使用prefect client删除flow后,还需要VACUUM表格吗
m
thought for 29.0 seconds
结论:需要,尤其是在你一次性删了大量 flow runs 的情况下。 原因与建议 - 用 Prefect Client 删除只是生成数据库里的 DELETE 操作,行会变成“死元组”。不做 VACUUM,空间不会真正回收,表和索引会膨胀,查询计划也可能变差。 - autovacuum 会在后台慢慢处理,但速度未必能跟上你的删除量,而且不会把文件物理变小,只是把空间留给后续重用。 - 最佳实践: - 大批删除后,手动执行 VACUUM (ANALYZE) 更新统计信息并清理死元组。 - 如果你需要“立刻”把磁盘占用降下来:用 pg_repack(更推荐,锁小)或在维护窗口用 VACUUM FULL(会持有较强锁)。 建议的最小维护命令(按你删除的对象为主)
Copy code
-- 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;
需要立即缩小文件体积(谨慎)
Copy code
-- 锁表强、影响写入;更推荐用 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 maintenance
r
@Marvin autovacuum难道不会自动清理吗
m
thought for 68.8 seconds
会,但有重要前提和局限: - autovacuum 会自动做 VACUUM/ANALYZE,清理“死元组”,让表内空间可被后续写入重用;也会更新统计信息,防止查询退化。 - 但它不会“缩小文件到操作系统”——也就是不会立刻降低磁盘占用。要想真正收缩文件体积,需要在维护窗口用 `VACUUM FULL`(强锁)或更推荐用 `pg_repack`(在线、锁更小)。 - 当你一次性删了大量行(比如用 Prefect Client 批量删除 flow runs)时,autovacuum 往往跟不上海量删除的速度,或者短期内不触发(阈值和比例是可配置的)。这时建议手动执行 `VACUUM (ANALYZE)`,必要时再做 `pg_repack`。 建议做法(结合你当前高吞吐场景) 1) 大删之后手动维护
Copy code
VACUUM (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 这类高写表 - 表级参数(不重启即可生效):
Copy code
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)示例:
Copy code
autovacuum = on
autovacuum_max_workers = 6
autovacuum_naptime = '10s'
autovacuum_vacuum_cost_limit = 2000
- 开启日志便于观察是否触发及时:
Copy code
ALTER SYSTEM SET log_autovacuum_min_duration = '0';
SELECT pg_reload_conf();
3) 监控是否“跟得上”
Copy code
-- 关注死元组和最近的自动/手动 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 maintenance
r
@Marvin 我自己搭建的prefect server每天处理flow有50w,prefect server部署多少实例合适
m
thought for 264.7 seconds
先给结论:Prefect Server(3.x 的 Orion API)是无状态的,通常通过“多副本 API + 独立的 PostgreSQL + 后台服务”来扩展。你每天 50 万个 flow run,平均约 5.8 个/秒(峰值可能更高)。没有一个“一刀切”的实例数,但可以用下面的基线配置起步,再按观测指标横向扩容。 建议的起步架构与实例数 - API(prefect server / Orion API) - 3–5 个副本起步,放在负载均衡后;每个副本 2 vCPU / 4 GB 内存左右 - 设置 HPA:当 CPU > 70% 或 p95 接口延迟 > 300 ms 时可扩到 6–8 个副本 - 注意数据库连接池总量(见下文) - 后台服务(scheduler、foreman、标记 late runs、事件持久化等) - 先跑 1 个副本;如果“计划运行创建延迟”(scheduler lag)> 60 秒,可扩到 2–3 个 - 这些服务是为并发而设计的,但扩太多只会增加数据库压力,建议小步渐进 - 数据库(PostgreSQL,最关键的瓶颈) - 独立托管或专用实例,建议 4–8 vCPU,16–32 GB 内存起步,按负载增长 - max_connections 至少 300+,推荐前置 PgBouncer 做连接池 - 关注 IOPS 与 Autovacuum;随着数据量增长尽早做表分区(flow_runs / task_runs 等大表) - 建议设置 plan_cache_mode=force_custom_plan,避免分区表上的泛化执行计划导致慢查询 为什么这个数量级通常足够 - 50w/天 ≈ 5.8 runs/s 平均。每个 flow run 会产生多次状态更新、任务运行、日志写入,数据库吞吐是核心瓶颈。 - 3–5 个 API 副本在数据库给力的情况下,通常可以覆盖该量级的平均吞吐;峰值或“整点风暴”时可加到 6–10 个。 关键参数与调优建议 - 数据库连接池(每个 API 副本) - 典型 SQLAlchemy 默认大约 pool_size≈5、max_overflow≈10,即每副本最多 ~15 个连接 - 粗略估算:5 个 API 副本 ≈ 75 个连接,再加后台服务和管理连接,Postgres 需留足余量 - 用 PgBouncer 做池化,可降低数据库压力 - Scheduler 调优(示例起步建议,按观测调整) - 将 insert_batch_size 提高到 500–1000 - 将 max_runs 提高到 200–500 - 将 loop_seconds 降到 15–30 秒 - 目标:降低“计划运行创建延迟”(从计划时间到被写入 DB 的延迟) - 日志与事件 - 如果日志量很大,考虑降低日志级别或开启批量写入/采样,减少对 DB 的写负载 - 负载平滑 - 使用全局并发限制(Global Concurrency Limits)避免瞬时洪峰 - 将大量定时(如整点)任务做抖动/错峰调度 扩容触发的观测指标(建议基于这些做 HPA/告警) - API - p95/p99 API 延迟、非 2xx 比例、每副本连接池饱和(等待时间/超时) - Scheduler - 从“计划时间”到“创建 run 记录”的延迟(lag) - 每分钟创建的 scheduled runs 速率是否跟得上 - 数据库 - CPU 利用率、活跃连接数与等待事件、慢查询比例、Autovacuum backlog - Workers/队列 - 队列 backlog、NOT_READY 状态、取工频率 需要你补充的信息(我们可以据此给出更精准的实例数) - 运行环境:Kubernetes 还是 VM?是否已有负载均衡/HPA? - Postgres 规格与是否有 PgBouncer? - 负载形态:均匀还是集中在整点/半点?峰值大致是多少倍平均? - 每个 flow 的任务数量/日志量级大概多少? - 是否将日志写到 API/DB,还是外部日志系统? 参考链接(强烈推荐) - 自托管参考(Docker Compose) - Prefect 设置参考大全 - 在生产 Postgres 上做分区(Prefect Cloud 实战) - Prefect Cloud 扩展中的连接回收与内存经验 - 全局并发限制(平滑洪峰) 总结 - 起步:API 3–5 个副本 + 后台服务 1 个副本 + 强壮的 Postgres(建议配 PgBouncer) - 若有明显峰值或调度落后,API 提升到 6–10,后台服务到 2–3,并调大 scheduler 批量参数 - 以数据库为中心做容量与观测,遇到 lag/慢查询时优先优化 DB(池化、参数、分区、索引)再加 API 副本 如果你能分享上面的环境与负载细节,我可以给出更具体的实例数建议和参数值。
r
@Marvin 给我一个可以定期清理flow/task表(最多保存15天)的完整的pgsql语句
m
thought for 223.6 seconds
下面给你两套“保留最近15天”的 PostgreSQL 清理脚本,适用于自托管的 Prefect 3.x(Orion)数据库。建议先在测试/备库验证,再在生产执行。 一、简单版(一次性事务,体量大时可能锁表) - 以 flow_run 的时间为准:COALESCE(end_time, created) < NOW() - interval '15 days' - 会按依赖顺序清理 task/flow 的状态、日志、工件等,再删 run 本体 - 如你的 schema/表名与默认不一致,请按需调整
Copy code
BEGIN;

-- 基准时间点(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 f
SET state_id = NULL WHERE f.id IN (SELECT id FROM to_delete_flow_runs); GET DIAGNOSTICS v_deleted = v_deleted + ROW_COUNT; -- 删除 task_run_state 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 ) DELETE FROM task_run_state s USING to_delete_task_runs r WHERE s.task_run_id = r.id; GET DIAGNOSTICS v_deleted = v_deleted + ROW_COUNT; -- 删除 flow_run_state WITH to_delete_flow_runs AS ( SELECT id FROM flow_run WHERE COALESCE(end_time, created) < cutoff_ts LIMIT batch_size ) DELETE FROM flow_run_state s USING to_delete_flow_runs r WHERE s.flow_run_id = r.id; GET DIAGNOSTICS v_deleted = v_deleted + ROW_COUNT; -- 删除日志 WITH to_delete_flow_runs AS ( SELECT id FROM flow_run WHERE COALESCE(end_time, created) < cutoff_ts LIMIT batch_size ), 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 ) DELETE FROM "log" l WHERE l.flow_run_id IN (SELECT id FROM to_delete_flow_runs) OR l.task_run_id IN (SELECT id FROM to_delete_task_runs); GET DIAGNOSTICS v_deleted = v_deleted + ROW_COUNT; -- 删除工件 WITH to_delete_flow_runs AS ( SELECT id FROM flow_run WHERE COALESCE(end_time, created) < cutoff_ts LIMIT batch_size ), 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 ) DELETE FROM artifact a WHERE a.flow_run_id IN (SELECT id FROM to_delete_flow_runs) OR a.task_run_id IN (SELECT id FROM to_delete_task_runs); GET DIAGNOSTICS v_deleted = v_deleted + ROW_COUNT; -- 删除 flow_run_input WITH to_delete_flow_runs AS ( SELECT id FROM flow_run WHERE COALESCE(end_time, created) < cutoff_ts LIMIT batch_size ) DELETE FROM flow_run_input i USING to_delete_flow_runs r WHERE i.flow_run_id = r.id; GET DIAGNOSTICS v_deleted = v_deleted + ROW_COUNT; -- 删除 task_run 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 ) DELETE FROM task_run t USING to_delete_task_runs r WHERE t.id = r.id; GET DIAGNOSTICS v_deleted = v_deleted + ROW_COUNT; -- 删除 flow_run WITH to_delete_flow_runs AS ( SELECT id FROM flow_run WHERE COALESCE(end_time, created) < cutoff_ts LIMIT batch_size ) DELETE FROM flow_run f USING to_delete_flow_runs r WHERE f.id = r.id; GET DIAGNOSTICS v_deleted = v_deleted + ROW_COUNT; -- 若本轮没有任何删除,跳出 IF v_deleted = 0 THEN EXIT; END IF; END LOOP; RAISE NOTICE 'Prefect cleanup finished'; END; $$; -- 调用示例(保留15天,每批5000) CALL prefect_cleanup(15, 5000);
Copy code
可选:用 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
@Marvin
Copy code
WITH 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有没有优化方法,执行很慢
m
thought for 56.3 seconds
瓶颈通常不在 SQL 写法本身,而是索引缺失和不必要的嵌套子查询。优化思路: - 确保有这3个关键索引(没有的话删除会非常慢) - 用 JOIN + USING 重写删除,避免大 IN 子查询 - 按“task_run 数量”分批,而不是按“flow_run 数量”分批 - 大删除后及时 ANALYZE/VACUUM,必要时用 pg_repack 回收空间 先建/补齐索引(并发建索引,避免长时间锁表)
Copy code
CREATE 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 分批)
Copy code
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) 如果你要进一步压低开销,可用临时表分两步(更清晰的执行计划、可重复使用)
Copy code
-- 第一步:预计算本批要删的 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 降低 - 每批删除后跑
Copy code
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
@Marvin 删除flow的sql脚本会不会有个问题:因为表中数据是按照顺序删除的,执行命令时会有时间差,最后删除flow_run导致部分task_run信息没有删除但flow_run已经被删了
m
thought for 74.3 seconds
不会出现“父被删了但子没删”的不一致,只要你的库里有外键约束: - 如果 task_run.flow_run_id 有外键但没有 ON DELETE CASCADE,Postgres 会阻止删除 flow_run(报错),你不可能把父删掉而留着子。 - 如果有 ON DELETE CASCADE,则删 flow_run 会自动级联删掉 task_run(以及其依赖,取决于各自的外键规则)。 但你说的“时间差”在分批脚本中确实会造成“本批挑选的集不同步”的问题(不是不一致,而是你删的集合前后不一致,可能留下一些子数据要下次再删)。解决办法是:在一个事务里“冻结本批ID集合”(用临时表一次性选出要删的 flow_run 和其 task_run),所有删除都基于这两个临时表执行。这样本批内无论顺序如何,都不会出现遗漏。 下面是一段稳妥的分批脚本(单事务、临时表冻结、顺序安全)。你可以在循环里重复调用它,或把 LIMIT 调大/小。若你的外键是非级联的,这段也能正常工作;若是级联的,你甚至可以省略部分子表删除。
Copy code
BEGIN;

-- 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:
Copy code
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
@Marvin 先解除可能的外键指针,避免约束冲突。这一步是否可以不用
m
thought for 105.0 seconds
是否可以不用,取决于你数据库里的外键“删除规则”。结论: - 可以不用的前提:相关外键是 ON DELETE CASCADE(或 ON DELETE SET NULL)并且你的删除顺序与之匹配。 - 否则不要省略。因为: - flow_run.state_id → flow_run_state、task_run.state_id → task_run_state 如果不是 ON DELETE SET NULL,你先删 state 会被外键拦下;反过来先删 run,如果 state 到 run 没有 CASCADE,也会被拦下。 - task_run.child_flow_run_id 和 parent_task_run_id 可能指向“本批不删除”的 flow_run/task_run,不置空会阻止删除。 先查一下你库里的外键删除规则
Copy code
-- 查关键外键的 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 过滤:
Copy code
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
@Marvin fk_table | fk_column | pk_table | delete_rule ----------------+-------------+----------------+------------- flow_run | state_id | flow_run_state | SET NULL flow_run_state | flow_run_id | flow_run | CASCADE task_run | flow_run_id | flow_run | CASCADE task_run_state | task_run_id | task_run | CASCADE
m
thought for 22.6 seconds
结论:根据你贴出的外键规则,可以省略“先置空外键指针”的步骤,直接删 flow_run 即可,Postgres 会自动级联清理相关的 task_run / flow_run_state / task_run_state。 你当前外键摘要 - flow_run.state_id → flow_run_state: ON DELETE SET NULL - flow_run_state.flow_run_id → flow_run: ON DELETE CASCADE - task_run.flow_run_id → flow_run: ON DELETE CASCADE - task_run_state.task_run_id → task_run: ON DELETE CASCADE 含义 - 删除 flow_run 时,会级联删除其所有 flow_run_state、task_run,进而再级联删除 task_run_state。 - flow_run.state_id 指向的状态被删时会自动 SET NULL;但因为 flow_run 自己也会被删,这里无需手动置空。 推荐的批量删除语句(只删 flow_run,其他表交给级联)
Copy code
-- 准备索引(只需建一次;没有这个索引会很慢)
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
@Marvin 根据以上我提供的信息,给我完整的sql方案,包括flow_run相关表清理和log表清理
m
thought for 79.7 seconds
基于你库里的外键规则: - flow_run.state_id → flow_run_state: ON DELETE SET NULL - flow_run_state.flow_run_id → flow_run: ON DELETE CASCADE - task_run.flow_run_id → flow_run: ON DELETE CASCADE - task_run_state.task_run_id → task_run: ON DELETE CASCADE 因此删除 flow_run 会自动级联清理 task_run / 两类 state;不需要手动置空/先删子表。你只需要: - 批量删除超过 15 天的 flow_run - 在删除 flow_run 前,显式清理引用这些 run 的日志(以及可选的工件、flow_run_input 等非级联表) 下面是一套完整方案(索引 + 批处理脚本)。你可以把批处理脚本用定时任务循环执行,直到返回 0 行为止。 一、建议的索引(只需执行一次)
Copy code
-- 按 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 会自动级联删除
Copy code
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 示例(可选)
Copy code
-- 如未启用
-- 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)。建索引后再按时间删除:
Copy code
-- 假设日志表时间列名为 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;
ANALYZE flow_run_state; ANALYZE task_run_state; ANALYZE "log"; ANALYZE artifact; ANALYZE flow_run_input; ``` - 空间回收交给 autovacuum;如需强制回收请用 pg_repack,避免 VACUUM FULL 造成长时间锁表。 - 执行前可先把 DELETE 改成 SELECT COUNT(*) 预估影响行数;或用 EXPLAIN (ANALYZE, BUFFERS) 检查执行计划与 I/O。 如果你还用到了其他与 run 相关的扩展表(事件、工件集合等),告诉我表名/主外键,我可以把对应的清理语句也补上。
r
@Marvin DELETE FROM "log" l这里的语句,为啥log要带引号
m
thought for 22.8 seconds
简短回答: - 双引号是 PostgreSQL 的“标识符引用”,用于精确指明表/列名。 - 对于 Prefect 默认的日志表名 log,在 PostgreSQL 中并不是保留字,且是全小写,所以一般不需要加引号;`DELETE FROM log l ...` 也能正常执行。 - 我在示例里加了引号是出于通用性/稳妥性(避免和某些工具/方言冲突)。如果你的表就是小写的 log,去掉引号完全没问题。 什么时候必须带引号 - 创建表时用了引号并包含大写或特殊字符(例如建的是 "Log" 或 "LOG"),之后所有引用都必须用双引号且大小写一致。 - 表名是保留字(log 在 PostgreSQL 里不是保留字,但在其他数据库/方言里可能是)。 快速自查表名到底是什么
Copy code
-- 两个都试一下,哪个不是 NULL 就说明实际注册的表名
SELECT to_regclass('public.log');
SELECT to_regclass('public."log"');

-- psql 里也可以
-- \d log
-- \d "log"
结论与建议 - Prefect 3.x 默认是小写 log 表,你可以直接用:
Copy code
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
@Marvin 你给我的sql脚本有个问题,删除log表数据这个命令执行非常慢
m
thought for 53.4 seconds
明白,删 log 慢通常有两个主要原因: - 使用 OR 和 IN 让优化器很难走到理想的索引路径 - 批量过大导致长事务/锁竞争/表膨胀 下面给你几步优化,按“影响最大 → 实施成本低”排序。你可以只替换“删除 log 的步骤”,其余 flow_run 的清理逻辑保持不变。 1) 先确保必要索引(只需建一次,注意 CONCURRENTLY 不能在事务里执行)
Copy code
CREATE 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) 替换你脚本里原先的
Copy code
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);
为下面两条(更容易命中索引):
Copy code
-- 先按 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):
Copy code
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;请确认你库里的列名)
Copy code
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
Copy code
ANALYZE log;
- 避免“长事务占着旧快照”拖慢 autovacuum:大批删除建议切成多次短事务 - 如果执行计划依然走顺序扫描,检查统计信息是否过期、索引是否被正确选择;必要时先 EXPLAIN (ANALYZE, BUFFERS) 对比“OR 写法”和“两条 JOIN 写法”的差异 把以上第 2 步替换进你原先的事务脚本通常就能显著提速;如果数据量特别大,再叠加第 3 步的 ctid 微批。需要我帮你把整段事务脚本改成“ctid 微批 + 多次提交”的版本,也可以直接贴出来。
r
@Marvin task抛出的异常raise RuntimeError(message=f"blade_req_final_failed={msg}"),怎么在flow中捕获并打印。我的业务场景中不知道是哪个task抛出的异常不能使用task返回值检查,只能使用在flow中使用tryexcept检查
m
thought for 278.4 seconds
假设你在用 Prefect 3.x。根据你如何调用 task,有两种常见方式在 flow 里“只用 try/except”捕获并打印 task 抛出的异常(比如 `RuntimeError(f"blade_req_final_failed={msg}")`)。 一、任务串行(直接调用 task) - 直接调用
my_task()
时,异常会像普通 Python 函数一样同步抛出;你可以在 flow 里用一个大的 try/except 包住整段逻辑。
Copy code
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 中捕获任意一个失败
Copy code
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) 想“全部执行完后一次性打印所有失败”的异常细节(不只第一个):
Copy code
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
@Marvin 以下是我的代码,帮我改造:@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 = 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 = {} tasks['b'] = b tasks['c'] = c if a.result()["value"] > 30: e = task_e.submit(a) tasks['e'] = e d = task_d.submit(tasks) 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()}") return { "flow_run_id": flow_run_id, "errno": 1, "errmsg": str(e) }
m
thought for 43.8 seconds
下面给你按 Prefect 3.x 改造了一版,满足“只在 flow 里用 try/except 捕获并打印,不逐个检查 task 返回值”的需求。要点: - 仍然用
submit()
并发提交;在 flow 顶层用 try/except 做总兜底 - 用
wait()
等所有已提交的 futures 完成,并通过
future.state
判断失败是谁,打印详细异常 - 只在需要时调用 `.result()`(比如你的条件判断和最终返回)
Copy code
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
@Marvin 我不要逐个检查task是否失败,这太麻烦了。
m
thought for 70.4 seconds
明白了——那就用“只在 Flow 里一次性取最终结果”的模式:不逐个检查 task,只在最后对聚合任务
d
调一次 `.result()`。如果任一上游 task(a/b/c/e/d 里任何一个)失败,Prefect 会在
d.result()
时抛出原始异常或 UpstreamTaskError,被你外层的 try/except 统一捕获并打印。 关键点:把对
a
的条件判断从 flow 移进
task_e
自身,这样在 flow 中完全不需要提前 `a.result()`。示例改造如下:
Copy code
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()`:
Copy code
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()
调用,自然就不存在“逐个检查”的负担。 可选等价写法(同样只需一次触发异常) - 用一次性解析触发异常:
Copy code
from prefect.futures import resolve_futures_to_results
resolve_futures_to_results([d])  # 任一上游失败会在这里抛出
result = d.result()
参考 - States 与上游失败传播 - PrefectFutures API
r
@Marvin UpstreamTaskError只能看到是哪个task run name失败了,不能获取到task raise出来的信息
m
thought for 290.2 seconds
确实,UpstreamTaskError 本身不包含上游 task 抛出的原始异常;要拿到“raise 出来的信息”,必须从失败 task 的 State 里提取异常。为了避免你在 flow 里“逐个检查”,给你一版“只在最终一次捕获,然后一次性把所有失败根因打印出来”的写法: - 仅在最后对
d.result()
做一次 try/except - 一旦捕获异常,就用
resolve_futures_to_states(...)
一次性把所有相关 futures 解析成 states - 从失败的 state 中用
get_state_exception(state)
提取真实异常并打印 - 不需要在主逻辑里逐个 if 检查哪个 task 失败 改造代码如下(延续我们之前把 e 的条件判断移进 task_e 的思路,避免在 flow 中提前 `a.result()`):
Copy code
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
@Marvin 讲一下怎么用on_failure来做吧
m
thought for 117.3 seconds
可以,用 on_failure 把“打印根因异常”的逻辑放到 task/flow 自己的失败钩子里,这样你在 flow 主体里就不需要逐个检查或者额外解析了。 核心要点 - Hook 触发时机:当 task/flow 最终进入 Failed 状态(重试耗尽后)才会触发。 - Hook 形参: - task 级别:`def hook(task, task_run, state): ...` - flow 级别:`def hook(flow, flow_run, state): ...` - 在 hook 里用
get_state_exception(state)
取到 task/flow 的真实异常;`state.message` 则是 Prefect 格式化后的错误信息。 - 日志:建议用 `get_run_logger()`,也可用你已有的 `mario_log`。 一、给 task 加 on_failure - 用装饰器或
with_options
都可以。你可以一个钩子复用在多个 task 上。
Copy code
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 级别钩子。
Copy code
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()
触发失败,不需要你手动逐个检查。
Copy code
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 hooks
r
@Marvin on_failure=[log_task_failure]怎么作为一个task的默认值,而不是每次定义一个新的task都写一遍
m
thought for 19.3 seconds
可以把 on_failure 钩子“默认化”,这样每次定义 task 不用重复写。常用有两种做法: 方案 A:在你自己的模块里重绑定 task(最简单) - 用 functools.partial 给
prefect.task
预设默认参数,然后在项目里都用你重绑定后的
task
装饰器。
Copy code
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,自动与默认钩子合并。
Copy code
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])
包装,然后项目里都用包装后的引用即可。 - 例如:
Copy code
from other.module import raw_task_b as _raw_task_b
    task_b = _raw_task_b.with_options(on_failure=[log_task_failure])
- flow 同理可以默认化(如果需要):
Copy code
from functools import partial
  from prefect import flow as _flow
  flow = partial(_flow, on_failure=[log_flow_failure])
文档 - State change hooks
r
@Marvin 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("on_failure", log_task_failure) return prefect_task(*args, **kwargs) 这种方式可以吗
m
thought for 21.3 seconds
可以,但你这版需要做两点修改: -
on_failure
需要是列表(可包含多个钩子),而不是单个函数; - 需要兼容两种用法:`@task` 和 `@task(...)`。因此要写成一个“既可直接装饰函数、也可带参数”的包装器,并且把默认钩子和用户自定义的钩子合并。 给你一份可直接用的“带默认 on_failure 的 task 包装器”,同时预设你提到的其它默认项:
Copy code
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)
用法
Copy code
@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
@Marvin def task(___fn_=None, **_kwargs_): """task""" # 统一设定部分参数默认值 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): return prefect_task(**kwargs)(fn) # 支持 @task 和 @task(...) return _decorate if __fn is None else _decorate(__fn) 我发现这个写法,在装饰类成员函数时无法支持
m
thought for 67.7 seconds
你这个包装在装饰“类的成员函数”时失效,原因是 Prefect 的
@task
返回的是一个 Task 对象,它不是描述符(不实现 `__get__`),所以作为类属性时不会自动绑定 `self`。写成
self.method.submit(...)
时,`self` 不会被自动作为第一个参数传入,导致行为异常。 给你三种可用方案,按简单程度排序: 方案 1(最简单):调用时用“类属性”而不是“实例属性” - 定义依然可以写在类里,但调用时用
ClassName.method.submit(self, ...)
显式传入 `self`。
Copy code
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。
Copy code
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)
用法:
Copy code
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
并把需要的上下文通过参数传入;或直接定义为模块级函数,最稳妥。 - 例如:
Copy code
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
@Marvin from prefect import task as prefect_task from prefect.cache_policies import NO_CACHE import logger.log as mario_log 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) # 你也可以用 mario_log mario_log.warning( f"act=task_run_failed task={task.name} " 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_): """task""" # 统一设定部分参数默认值 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): t = prefect_task(**kwargs)(fn) return _MethodFriendlyDescriptor(t) # 支持 @task 和 @task(...) return _decorate if __fn is None else _decorate(__fn) 我的代码是这样,用了发现不行,类成员task没有给到默认的on_failure
m
thought for 28.6 seconds
你这段代码里有不少变量名混用的问题(例如参数叫 _state_,但用的是 state;_task_obj_,但用的是 task_obj;等等),导致装饰成功但 on_failure 没真的挂上去。下面给你一份“可直接用”的版本,满足: - 统一默认项:cache_result_in_memory=False、persist_result=False、cache_policy=NO_CACHE、on_failure=[log_task_failure] - 支持 @task 和 @task(...) 两种写法 - 支持装饰类成员函数,实例上调用 self.method.submit(...) 会自动注入 self - 也支持 ClassName.method.submit(self, ...) 的传统用法 - on_failure 可被追加(不会覆盖默认钩子) 直接粘贴使用:
Copy code
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)}"
    )
    # 如果要同步写你们的日志系统:
    # 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)
用法示例(类成员函数)
Copy code
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 hooks
r
@Marvin 我的类成员task是这样是这样使用的,用起来不行:blades_client = BladesTask() blades_req = blades_client.req blades_req.submit(xxxx)
m
thought for 48.6 seconds
明白,你想用实例方式写成: - blades_client = BladesTask() - blades_req = blades_client.req - blades_req.submit(xxx) 这种写法可以实现,但前提是你的 @task 装饰器是“方法友好”的版本(带描述符 + 代理),且不要被实例属性遮蔽。你之前的代码里有不少变量名混用/遮蔽问题,导致最终拿到的不是代理对象,从而 on_failure 也没挂上。 给你一份可直接用的最小可运行版本,满足: - @task 和 @task(...) 两种写法 - 类成员方法支持 self.method.submit(...),自动注入 self - 默认带 on_failure=[log_task_failure],也可叠加用户自定义 hooks - 不会被 with_options 打破 包装器代码(放到你的公共模块里,比如 my_prefect_wrappers.py)
Copy code
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)}"
    )
    # 如需你们的日志: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)
使用示例(符合你的写法)
Copy code
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(...)`。 - 想确认是否生效:打印类型
Copy code
print(type(BladesTask.req), type(blades_client.req))
  # 应该显示 _MethodFriendlyDescriptor / _BoundTaskProxy(或内部包装类型)
如果你把这段放进项目后依旧不行,麻烦贴一下具体的报错信息(TypeError/AttributeError 等)和你类中
req
的定义(是否在 init 里重名覆盖),我再帮你对症修正。
r
@Marvin 我打印显示的是<class 'prefect.tasks.Task'> <class 'workflow.utils.task._BoundTaskProxy'>
m
thought for 77.0 seconds
这俩类型是对的: - BladesTask.req → <class 'prefect.tasks.Task'>(类属性,用于 Class.req.submit(self, ...)) - blades_client.req → <class '..._BoundTaskProxy'>(实例属性,支持 self.req.submit(...) 自动注入 self) 如果 on_failure 没触发,通常是以下原因: - 任务最终没有进入 Failed(例如有重试并最终成功)——on_failure 只在“最终失败”时触发 - 某处没有用你的自定义 task 包装器(用了原生 prefect.task) - 用了 ClassName.req.with_options(on_failure=[...]) 覆盖了默认钩子 - 在 init 里写了 self.req = ...,把类属性描述符遮蔽了 给你一版“可直接用”的包装器,修正了你之前的变量名问题,并且让实例上的 with_options 在用户传 on_failure 时,自动把“默认钩子 + 用户钩子”合并(避免把默认钩子覆盖掉):
Copy code
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)}"
    )
    # 如需你们的日志系统:
    # 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 hooks
r
@Marvin 另一种方案:class Svc: @task # 你的带默认 on_failure 的包装 def step(self, x): ... svc = Svc() # 不要用 svc.step.submit(x) # 要这样: Svc.step.submit(svc, x) # 显式把 self 作为第一个参数传进去,能不能优化一下,定义一个blades_req,使用blades_req.submit调用task
m
thought for 22.7 seconds
可以优化,给你两种不改 Prefect 本体的写法,满足“定义一个 blades_req,再用 blades_req.submit(...)”: 方案 1:用我们之前的“方法友好版 @task”(带描述符) - 已经用过那个包装的话,直接:
Copy code
blades_client = BladesTask()
blades_req = blades_client.req   # 这是代理对象
blades_req.submit(xxx)
这就是你要的形式,on_failure 默认钩子也会生效。 方案 2:保持原生 Prefect @task,写一个轻量“绑定器”把类方法 Task 绑定到实例 - 适用于 Svc.step 是 <class 'prefect.tasks.Task'> 的场景。 先写个通用绑定器:
Copy code
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)
然后这样用:
Copy code
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)
可选:做成类里的属性,使用更顺手
Copy code
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
@Marvin 这两种方法我用了都不行,帮我看看是啥问题,我的类成员定义方式是:class BladesTask(): """BladesTask """ def __init__(self, bns=blades_bns_online): """ bns_cache: 提供 get_bns(bns) 方法的服务发现缓存 prefix: BNS 前缀(默认 "bns://") """ 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, _video_url_=None, _text_utf8_=None, _feat_args_=None, _ip_port_=None, _retry_max_=3, _timeout_s_=60, _cos_name_="req", _log_id_: int=0, _output_format_="auto", _save_file_=False): """ blades_req """ raise RuntimeError(f"blade_req_final_failed=1")
m
thought for 65.2 seconds
看到了,你现在的类型打印是 - BladesTask.req → <class 'prefect.tasks.Task'>(类属性) - blades_client.req → <class '..._BoundTaskProxy'>(实例属性) 这两个是对的,说明“方法友好版”的包装器的描述符已经生效。现在“用起来不行”的根因主要有三类问题,你的代码里至少命中了前两个: 1) on_failure 钩子内部变量名用错 - 你定义的是 `def log_task_failure(task, _task_run_, state)`,但内部用的是
state
和 `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):
Copy code
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,
video_url=None, text_utf8=None, feat_args=None, ip_port=None, retry_max=3, timeout_s=60, cos_name="req", log_id: int = 0, output_format="auto", save_file=False, ): # 故意失败,测试 on_failure raise RuntimeError("blade_req_final_failed=1")
Copy code
使用方式(你想要的 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
@Marvin 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() class BladesTask(): """BladesTask """ def __init__(self, bns=blades_bns_online): """ bns_cache: 提供 get_bns(bns) 方法的服务发现缓存 prefix: BNS 前缀(默认 "bns://") """ 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, _video_url_=None, _text_utf8_=None, _feat_args_=None, _ip_port_=None, _retry_max_=3, _timeout_s_=60, _cos_name_="req", _log_id_: int=0, _output_format_="auto", _save_file_=False): """ blades_req """ upload_log = get_run_logger() upload_log.info(f"local_ip={get_local_ip()}") blade_info = f"feat_id={feat_id}, user_id={user_id}, token={token}, cos_name={cos_name}, log_id={log_id}" upload_log.info(f"{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 = {} req_obj["user"] = {"user_id": user_id, "token": token, "group_id": group_id} req_obj["log_id"] = log_id req_obj["feat_args"] = [{"feat_id": feat_id, "feat_args": feat_args}] req_obj["input_data"] = { "video_url": video_url, "text_utf8": text_utf8_value } req_str = json.dumps(req_obj).encode(encoding='UTF8') if save_file: save_obj = text_utf8 if save_obj is None: save_obj = req_str fs.save(query=log_id, _file_name_=f"blades-{feat_id}-{cos_name}_input", text=save_obj) # 请求 result = None msg = "" async with httpx.AsyncClient(timeout=timeout_s) as client: for retry_time in range(retry_max): await asyncio.sleep(0.5 + 2 * retry_time) try: if ip_port is None: ip, port = self._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}") mario_log.notice(f"act=blade_req_post, {blade_info}, retry_time={retry_time}") # 异步请求 response = await 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"] # 失败 if feat_status != 0: msg = resp_json['feat_res'][0]['feat_detail_status'] upload_log.warning(f"feat_detail_status={feat_status}, err_msg={msg}") mario_log.warning(f"act=blade_req_failed, {blade_info}, " f"retry_time={retry_time}, msg={msg}") # 判断是否重试 if feat_status >= 5002 and feat_status <= 5500: # 网络错误与可重试错误 continue elif feat_status >= 5600 and feat_status < 5799: # 限流错误sleep后重试 continue else: break # 成功 result = base64.b64decode(resp_json["feat_res"][0]["value"]) break except Exception as e: upload_log.warning(f"exception={str(e)}") mario_log.warning(f"act=blade_req_raise, {blade_info}, " f"traceback={traceback.format_exc()}, msg={msg}") del text_utf8_value del req_obj del req_str del resp_json # 请求失败 if result is None: mario_log.warning(f"act=blade_req_final_failed, {blade_info}, " f"msg={msg}") raise RuntimeError(f"blade_req_final_failed, feat_id={feat_id}, msg={msg}") # 请求成功 # 转换 bytes to str if isinstance(result, bytes): result = result.decode("utf-8") result = result.strip() result_len = len(result) upload_log.info(f"blade_req_success, result_len={result_len}") # 自动判断返回类型 if output_format == "auto": if result.startswith('{') and result.endswith('}'): result = json.loads(result) if save_file: fs.save(query=log_id, _file_name_=f"blades-{feat_id}-{cos_name}_result", text=result) return result 我想把这个类成员函数task改造成普通函数task,怎么弄
m
thought for 69.0 seconds
看起来你现在把
@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 视为失败
if feat_status != 0: msg = resp_json["feat_res"][0].get("feat_detail_status", "") upload_log.warning(f"feat_detail_status={feat_status}, err_msg={msg}") if mario_log: mario_log.warning( f"act=blade_req_failed, {blade_info}, retry_time={retry_time}, msg={msg}" ) # 可重试错误 if 5002 <= feat_status <= 5500 or 5600 <= feat_status < 5799: ip_port = None # 下次重试换一个 host continue else: break # 成功 result = base64.b64decode(resp_json["feat_res"][0]["value"]) break except Exception as e: upload_log.warning(f"exception={str(e)}") if mario_log: mario_log.warning( f"act=blade_req_raise, {blade_info}, traceback={traceback.format_exc()}, msg={msg}" ) ip_port = None # 下次重试换 host # 请求失败 if result is None: if mario_log: mario_log.warning(f"act=blade_req_final_failed, {blade_info}, msg={msg}") raise RuntimeError(f"blade_req_final_failed, feat_id={feat_id}, msg={msg}") # bytes -> str if isinstance(result, bytes): result = result.decode("utf-8") result = result.strip() result_len = len(result) upload_log.info(f"blade_req_success, result_len={result_len}") # 自动判断返回类型 if output_format == "auto": if isinstance(result, str) and result.startswith("{") and result.endswith("}"): result = json.loads(result) if save_file and fs is not None: fs.save(query=log_id, file_name=f"blades-{feat_id}-{cos_name}_result", text=result) return result
Copy code
如何在 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
Copy code
说明
- 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
@Marvin def log_task_failure(task, _task_run_, state): """log_task_failure """ # 这里不会重试,说明已是最终失败 exc = get_state_exception(state) # 你也可以用 mario_log mario_log.warning( f"act=task_run_failed task={task.name} " f"exc={exc}\n{exception_traceback(exc)}") def task(___fn_=None, **_kwargs_): """task""" # 统一设定部分参数默认值 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): return prefect_task(**kwargs)(fn) # 支持 @task 和 @task(...) return _decorate if __fn is None else _decorate(__fn) 我发现这个写法不支持async函数task
m
thought for 76.2 seconds
你这个包装器“不支持 async”的根因不是 Prefect 限制,而是装饰器里有几处变量名/调用错误,导致任何函数(更别说 async)都包不对: -
_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(...)
两种用法都支持):
Copy code
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
使用示例
Copy code
@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)
Copy code
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
@Marvin def log_task_failure(task, task_run, state): """log_task_failure """ # 这里不会重试,说明已是最终失败 exc = get_state_exception(state) # 你也可以用 mario_log mario_log.warning( f"act=task_run_failed task={task.name} " f"exc={exc}\n{exception_traceback(exc)}") def task(__fn=None, **kwargs): """task""" # 统一设定部分参数默认值 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): wrapped = prefect_task(**kwargs)(fn) return wrapped # 支持 @task 和 @task(...) return _decorate if __fn is None else _decorate(__fn) 我发现这个写法不支持async函数task
m
thought for 58.2 seconds
你这版包装器本身是可以支持 async 的,Prefect 3.x 的
task
装饰器对异步函数天然支持。通常出现“好像不支持 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:
Copy code
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
Copy code
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 执行并返回结果)
Copy code
from prefect import flow

@flow
def sflow():
    y = my_async_task(1)  # 不需要 await
    return y
- Flow 外(纯 Python 环境):async 任务需要 await 或 asyncio.run
Copy code
# 在交互式/脚本里这么用
import asyncio
asyncio.run(my_async_task(1))
如果你还是遇到“不支持 async”,请贴一下调用场景(在 Flow 里还是 Flow 外)和完整报错(尤其是是否有 “coroutine was never awaited”)。我可以据此进一步定位。文档参考:Prefect 3.x Tasks
r
@Marvin Traceback (most recent call last): File "/home/gelin/gitcode/baidu/miaobi/baikan-video-workflow/env-output/bin/des/python_env/self_site_packages/cpu/prefect/task_engine.py", line 950, in call_hooks result = hook(task, task_run, state) ^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/home/gelin/gitcode/baidu/miaobi/baikan-video-workflow/env-output/bin/workflow/utils/task.py", line 29, in log_task_failure f"exc={exc}\n{exception_traceback(exc)}") ^^^^^^^^^^^^^^^^^^^^^^^^ File "/home/gelin/gitcode/baidu/miaobi/baikan-video-workflow/env-output/bin/des/python_env/self_site_packages/cpu/prefect/exceptions.py", lin e 52, in exception_traceback tb = traceback.TracebackException.from_exception(exc) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/home/gelin/gitcode/baidu/miaobi/baikan-video-workflow/env-output/bin/des/python_env/lib/python3.11/traceback.py", line 830, in from_ex ception return cls(type(exc), exc, exc.traceback, *args, **kwargs) ^^^^^^^^^^^^^^^^^ AttributeError: 'coroutine' object has no attribute '__traceback__' log_task_failure有这个错误怎么解决,可能因为task是async的
m
thought for 47.3 seconds
问题根因:`get_state_exception(state)` 在你当前环境里返回了一个 coroutine,你把它直接传给了 `exception_traceback()`,于是报了 “‘coroutine’ object has no attribute __traceback__”。这通常发生在: - 导入到了异步版的异常提取函数(或跨版本差异),导致需要 await 才能拿到真正的 Exception; - 但 Prefect 的
on_failure
hook 是同步调用的(引擎没有 `await`),所以在 hook 里不能直接 `await`。 推荐两种安全修法(均为同步,兼容 async/sync task): 方案 A:直接使用 Prefect 提供的格式化函数(最简洁,零 await)
Copy code
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;拿不到就回退到字符串
Copy code
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
@Marvin def log_task_failure(task, _task_run_, state):函数里能否获取到task的run_name
m
thought for 29.8 seconds
可以的。在 Prefect 3.x 的 task on_failure hook 里,第二个参数
task_run
就是当前任务运行实例,你可以直接拿到运行名: - 任务定义时设置的 `task_run_name`(包括用参数格式化后的最终名字),在运行时会被解析并存到
task_run.name
-
task.name
是任务定义名(非运行名) 示例:
Copy code
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 里拿到的是替换后的最终名字):
Copy code
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
@Marvin def log_task_failure(task, _task_run_, state): 里怎么让async函数的异常打印出来完整堆栈?
m
thought for 58.4 seconds
可以。建议用 Prefect 自带的格式化函数,它会返回包含完整堆栈的字符串,且不依赖你去“await”异常(on_failure hook 是同步调用的)。 最推荐(简单可靠,包含完整堆栈)
Copy code
from 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,否则回退”的写法:
Copy code
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
@Marvin format_exception里没有类似于traceback样式的错误信息,比如异常发生在哪一行,我需要这些信息
m
thought for 134.5 seconds
可以。思路是:在 on_failure 钩子里把真正的 Exception 对象拿出来,再用
prefect.exceptions.exception_traceback(exc)
生成完整堆栈(含行号)。困难在于某些环境下
get_state_exception(state)
返回的是 coroutine;on_failure 是同步调用,不能直接 await。可以用 anyio 的桥接在同步钩子里“间接 await”。 给你一份稳健实现(优先拿 Exception;拿不到就兜底用 `format_exception`):
Copy code
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
@Marvin ImportError: cannot import name 'aget_state_exception' from 'prefect.states'
m
thought for 51.0 seconds
出现 ImportError 说明你用的 Prefect 版本里没有 `aget_state_exception`。没关系,你仍然可以在同步的 on_failure 钩子里拿到“完整 traceback(含文件与行号)”,方法是用同步的
raise_state_exception(state)
把异常抛出来,再在钩子里捕获并格式化。 给你一版稳健实现(不依赖异步,兼容你当前版本):
Copy code
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
@Marvin我的版本是3.4.12
m
thought for 30.0 seconds
收到。Prefect 3.4.12 且没有 aget_state_exception 时,推荐在 on_failure 钩子里用同步的 raise_state_exception(state) 把真实异常抛出再捕获,这样能拿到完整 traceback(含文件与行号),同时兼容 async/sync 任务。 可直接用下面这版(保持你喜欢的参数名风格),纯同步、稳健:
Copy code
import 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
@Marvin 132611.271 | INFO | Task run 'blades-9039-text' - Running hook 'log_task_failure' in response to entering state 'Failed' /home/gelin/gitcode/baidu/miaobi/baikan-video-workflow/env-output/bin/workflow/utils/task.py35 RuntimeWarning: coroutine 'sync_compatible.<lo cals>.coroutine_wrapper.<locals>.ctx_call' was never awaited raise_state_exception(state) # 这里会 raise 真实异常 RuntimeWarning: Enable tracemalloc to get the object allocation traceback
m
thought for 60.7 seconds
你遇到的 “coroutine was never awaited” 是因为
raise_state_exception(state)
在你这个版本/上下文里返回了一个 coroutine,但 on_failure 钩子是同步调用的,Prefect 不会帮你 `await`。解决办法:在钩子里用 anyio 的线程桥接把这个 coroutine 跑起来,再捕获真实异常并格式化 traceback(含文件/行号)。 给你一份稳健可用的实现(Prefect 3.4.12,适配 async/sync 任务,钩子保持同步):
Copy code
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
@Marvin run_flow_failed_arise=Failed to resolve inputs in parameter 'server_record'. If your parameter type is not supported, consider using the
quote
annotation to skip resolution of inputs. 这是什么错误?
m
thought for 284.8 seconds
这是 Prefect 在“解析任务/Flow 入参”阶段报的错:引擎会尝试对每个参数做解析(递归遍历集合、把 Prefect Future/占位符解引用、做映射等)。如果某个参数的类型不受支持(比如自定义对象、打开的连接/句柄、协程对象、复杂不可序列化的实例),解析器会失败并抛出这条提示。对这种值,应该告诉 Prefect “不要解析,按原样传入”——用 quote 注解即可。 两种常见修法(Prefect 3.4.12): - 在参数上做类型注解(推荐,省心)
Copy code
from 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)  # 这里直接传即可
- 在调用处对值做包装
Copy code
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
@Marvin Failed to resolve inputs in parameter 'upload_data'. If your parameter type is not supported, consider using the
quote
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,看打印结果确实是json
m
thought for 295.6 seconds
这个错误的意思是:在执行
upload_to_bos
任务之前,Prefect 会递归“解析(resolve)”你传入的参数,把里面的 Prefect futures、可迭代容器等都遍历并转成最终值。如果在遍历
upload_data
这个参数时遇到 Prefect 不支持/无法遍历的类型,就会报 “Failed to resolve inputs in parameter 'upload_data'… consider using the
quote
annotation”。 结合你这段代码:
Copy code
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) 示例:
Copy code
@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,然后返回。 示例:
Copy code
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 的字段:
Copy code
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_inputsvisit_collection 如果你能贴一下
record_data
的结构或每个字段的类型(例如是否包含 pandas/numpy/自定义类/bytes),我可以更精准地建议是转型还是在哪些字段上使用 `quote`。
r
@Marvin {"input_json": "dict", "text_res": "dict", "planning_tts_json": "dict", "tts_res": "dict", "ttv_res": "PrefectFutureList", "video_path_list_url": "dict", "cover_res": "dict", "card_res_list": "PrefectFutureList", "aigc_video_res_list": "PrefectFutureList", "voice_effect_res": "dict", "video_effect_res": "dict"} 这是我打印出来的类型,有问题吗?
m
thought for 72.3 seconds
短答:你这份“参数类型”里出现 PrefectFutureList 本身不一定是问题;真正报错点在于下一个任务
upload_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 等)
Copy code
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 有问题的字段)
Copy code
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_inputsvisit_collection
r
@Marvin @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_): """preproc_upload_to_bos """ ret = { "input_json": input_json, "text_res": text_res, "planning_tts_json": planning_tts_json, "tts_res": tts_res, "ttv_res": ttv_res, "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()} fs.save(query=input_json["task_id"], _file_name_=f"preproc_upload_to_bos", text=ret) fs.save(query=input_json["task_id"], _file_name_=f"preproc_upload_to_bos_type", text=ret_type) return ret 这是我打印类型的代码,有问题吗
m
thought for 49.4 seconds
有关键问题:你在 task 的参数名是带下划线的(例如 `_input_json_`),但在函数体里用的却是没有下划线的名字(例如 `input_json`)。这会导致你实际用到的是“全局作用域里同名变量”(很可能还是 PrefectFuture/PrefectFutureList),而不是 Prefect 已经解析好的入参。结果就是你把“未解析的 futures 或复杂对象”塞进了返回的 `ret`,下游
upload_to_bos(upload_data=record_data)
解析
upload_data
时就会失败,报你看到的 “Failed to resolve inputs in parameter 'upload_data'…”。 修复方法(首选) - 在 task 里用“参数本身”的名字,不要用全局变量。最好把参数名改成一致的、无下划线的名字,函数体里就用同名变量。 示例(修正后)
Copy code
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)
Copy code
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
@Marvin @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): """preproc_upload_to_bos """ ret = { "input_json": input_json, "text_res": text_res, "planning_tts_json": planning_tts_json, "tts_res": tts_res, "ttv_res": ttv_res, "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()} fs.save(query=input_json["task_id"], file_name=f"preproc_upload_to_bos", text=ret) fs.save(query=input_json["task_id"], file_name=f"preproc_upload_to_bos_type", text=ret_type) return ret 其中,ret_type打印的结果为:{"input_json": "dict", "text_res": "dict", "planning_tts_json": "dict", "tts_res": "dict", "ttv_res": "PrefectFutureList", "video_path_list_url": "dict", "cover_res": "dict", "card_res_list": "PrefectFutureList", "aigc_video_res_list": "PrefectFutureList", "voice_effect_res": "dict", "video_effect_res": "dict"} 能看出问题吗?
m
thought for 333.4 seconds
能看出来问题:在你的 task 里,`ttv_res`、`card_res_list`、`aigc_video_res_list` 仍然是 `PrefectFutureList`。如果这个 task 是在一个 Prefect 3.x 的 flow 里被正常调度执行,参数会在进入 task 前被“解析”为最终值,按理这里应该是普通的 Python `list`,而不是 `PrefectFutureList`。这也是下游
upload_to_bos(upload_data=record_data)
解析
upload_data
时失败的根因。 这通常有两种原因: - 该 task 并没有在一个
@flow
运行上下文中执行(比如直接在普通 Python 代码里调用)。在这种情况下,参数不会被 Prefect 解析。 - 上游把一个
PrefectFutureList
直接当值往下传,且解析阶段没有被展开(例如某些场景/封装导致 Prefect 没有递归进去)。 快速修复方案(任选其一或组合): 1) 确保是在 flow 里调用 示例:
Copy code
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 列表解析成结果列表
Copy code
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`):
Copy code
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
@Marvin BEGIN; SELECT COUNT(*) FROM flow_run WHERE end_time IS NOT NULL AND end_time < now() - interval '12 days'; -- 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 '12 days' ORDER BY id LIMIT 10000; -- 调整批量大小以平衡时长与锁 -- 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 USING _del_flow_runs d WHERE l.flow_run_id = d.id; DELETE FROM log l USING _del_task_runs t WHERE l.task_run_id = t.id; -- 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; 这是我正在使用的例行清理pgsql数据的脚本,每天需要清理50w flow,但是每次运行太慢了,尤其是log表清理很慢,有没有更快的方式,比如log表直接根据时间清理
m
thought for 230.5 seconds
这类清理的瓶颈几乎都在
log
表:行数巨大、删除会产生大量 dead tuples、需要频繁 VACUUM;再加上没有良好索引或按时间清理就更慢。给你几个层次的优化方案,按“立刻可用 → 架构级提速”排序。 快速提速(尽量不改架构) - 先按时间直接清理日志,再清理 flow_run。对于“保留 12 天”的策略,日志和 flow_run 的生存期一致,优先按
timestamp
清理日志就能避免你现在的两次 JOIN 删除。 - 给日志相关列建合适的索引,特别是时间列用 BRIN,ID 列用 BTree:
Copy code
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):
Copy code
-- 每次删 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 SettingsLogging 定制 时间分区(架构级,“删除变成丢分区”) - 把
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 可以去掉或极大缩小批量):
Copy code
-- 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 APIPython models.logs - 数据库维护指南(说明事件保留、日志表等):Database maintenance - 日志定制:Logging customization 能否补充几条信息,我可以给出更精确的 DDL/迁移脚本建议: - 你使用的 Postgres 版本? - 是否可以接受把
log
表改为按时间分区? - 是否可以降低日志保留期,或者降低/关闭把日志写入 Prefect API(UI 中就看不到历史日志了)?
r
@Marvin task失败的hook函数中(def log_task_failure(task, _task_run_, state)),能否获取到flow的信息,比如flow_run_name
m
thought for 143.7 seconds
可以的。在 Prefect 3.x 的任务失败 hook 里可以拿到父 Flow 的运行信息。最直接的方式是通过
FlowRunContext.get()
读取当前的 Flow 上下文,然后访问其中的 `flow_run.name`(也就是 `flow_run_name`)。 示例(同步 hook):
Copy code
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:
Copy code
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
@Marvin prefect3 我写了一个flow代码,怎么在flow执行失败时或缺到失败的task:@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 = 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 = {} tasks['b'] = b tasks['c'] = c if a.result()["value"] > 30: e = task_e.submit(a) tasks['e'] = e d = task_d.submit(tasks) return { "flow_run_id": flow_run_id, "errno": 0, "errmsg": "", "result": d.result() } except Exception as e: return { "flow_run_id": flow_run_id, "errno": 1, "errmsg": str(e) }
m
thought for 598.9 seconds
可以拿到。常见有两种做法,给你对应的代码模板: 方式 A:在 flow 里用 futures 收集失败的任务(不等 flow 整体失败也能拿到)
Copy code
from 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}")
@flow(on_failure=[list_failed_tasks_on_flow_failure], flow_run_name="{logid}") def my_flow(input_data, logid): ... ``` 补充: - 获取 `flow_run_id`/`flow_run_name`:`from prefect.context import FlowRunContext; ctx = FlowRunContext.get(); flow_run_id = ctx.flow_run.id; flow_run_name = ctx.flow_run.name` - 如果某任务因为上游失败被跳过/取消,状态通常是
CANCELLED
或 `CRASHED`,可以一并统计。 相关源码参考: - Futures/状态 API:prefect.futures 源码 - 客户端过滤器(TaskRunFilter 等):filters.py - Orchestration 客户端:orchestration client 如果你愿意,我可以根据你现有的
task_a/b/c/d/e
具体返回结构,帮你把失败聚合的返回结构再细化一下。
r
@Marvin failed_tasks_name=[] local_vars = locals() for var_name, var_val in local_vars.items(): if isinstance(var_val, PrefectFuture): state = var_val.get_state() if state.is_failed() or state.is_crashed(): failed_tasks_name.append(var_name) elif isinstance(var_val, (list, PrefectFutureList)): if len(var_val) == 0: continue first_item = var_val[0] if isinstance(first_item, PrefectFuture): for index, fut in enumerate(var_val): state = fut.get_state() if state.is_failed() or state.is_crashed(): failed_tasks_name.append(var_name) break 这样判断失败的task可以吗
m
thought for 23.5 seconds
思路基本可行,但有几点需要调整,才能在 Prefect 3.x 下稳定工作: 关键修正 - PrefectFuture 没有
get_state()
方法,直接用
future.state
读取状态。 - 建议先
wait()
一下,确保所有 future 都已完成,否则你读取到的可能是 PENDING/RUNNING。 - 如果你把“失败”定义得宽一点,通常会把 CRASHED/CANCELLED 也算进来(很多时候下游因上游失败而被取消)。 修正版(保留你用 locals() 扫描的写法)
Copy code
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。
Copy code
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
@Marvin 我发现format_exception获取到的错误信息比较冗余,能不能获取简介的错误。比如task中抛出raise RuntimeError(f"blade_req_final_failed"),我想只获得blade_req_final_failed这个字符串
m
thought for 43.5 seconds
可以。不要用 `format_exception`,直接拿到异常对象再取它的 message 就行,最精简。 推荐两种安全方式(都不会再次抛错): 方式 1:从 future 直接取异常对象
Copy code
# 假设 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:从状态取异常对象
Copy code
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 里同理:
Copy code
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”,如果你只想冒号后的部分,也可以做个简单切分:
Copy code
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.states
r
@Marvin def get_failed_tasks(_local_vars_): """ 获取失败的任务名称 """ failed_task_name = "" failed_task_msg = "" for var_name, var_val in local_vars.items(): # 检查 PrefectFuture if isinstance(var_val, PrefectFuture): if var_val.state.is_failed(): failed_task_name = var_name failed_task_msg = get_err_msg(var_val.state) break # 检查 PrefectFutureList elif isinstance(var_val, (list, PrefectFutureList)): if len(var_val) == 0: continue first_item = var_val[0] if isinstance(first_item, PrefectFuture): for index, fut in enumerate(var_val): if fut.state.is_failed(): failed_task_name = var_name failed_task_msg = get_err_msg(var_val.state) break return failed_task_name, failed_task_msg 这段获取失败task name的函数有问题,failed_task_name写的是变量名,帮我改成task名
m
thought for 77.1 seconds
可以,把“变量名”改成“任务名”(更准确地说是 Task Run 名称)需要通过 future 的
task_run_id
去读对应的 TaskRun,再取 `TaskRun.name`。下面是修正后的函数,也顺手修了你原来 list 分支里取
var_val.state
的小 bug,以及外层/内层循环的 break 逻辑。 说明: - PrefectFuture 本身没有
task_name
属性;要么查询 Orchestration API 的 `TaskRun.name`(运行名),要么在提交任务时自己维护一个“任务定义名 → future”的映射。 - 我这里返回的是 Task Run 名称(通常就是你 Task 的可读名),以及精简的错误信息。 版本 1:直接查 TaskRun(最简单,零改动现有提交流程)
Copy code
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),直接用你希望的任务名键。
Copy code
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 "", ""
示例用法:
Copy code
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 client
r
@Marvin File "/home/gelin/gitcode/baidu/miaobi/baikan-video-workflow/env-output/bin/workflow/utils/flow_task_helpler.py", line 37, in get_failed_tasks task_run = client.read_task_run(var_val.task_run_id) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/home/gelin/gitcode/baidu/miaobi/baikan-video-workflow/env-output/bin/des/python_env/self_site_packages/cpu/prefect/client/orchestration/__init__.py", line 1509, in read_task_run raise prefect.exceptions.ObjectNotFound(http_exc=e) from e prefect.exceptions.ObjectNotFound: None 方式查询 Orchestration API好像不行
m
thought for 57.0 seconds
你这个报错说明通过 Orchestration API 查不到该 task run(ObjectNotFound)。常见原因: - 运行未连接到同一个 Prefect API(本地离线执行或 profile/workspace 配置不一致) - 正在用的
get_client
指向的 API 与当前 flow run 的 API 不同 - 极少数情况下,ID 不是 task run(比如 subflow 的 flow run future) 为了稳妥拿到“任务名”,我给你两套可用方案: 方案 A(推荐):提交时维护“任务名 → future”的映射,不依赖 API - 优点:不需要调用 API,离线/本地都稳定 - 你可以用
task_obj.name
作为键,这就是任务定义名(接近你想要的“task名”)
Copy code
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 "", ""
用法示例:
Copy code
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, msg
# PrefectFutureList 或 list[PrefectFuture] elif isinstance(var_val, (list, PrefectFutureList)): for fut in var_val: if not isinstance(fut, PrefectFuture): continue st = fut.state if st.is_failed() or st.type in {StateType.CRASHED, StateType.CANCELLED}: msg = concise_err(st) name = name_from_api(fut) or name_from_map(fut) or var_name return name, msg return "", "" ``` 建议 - 更稳的做法是“提交时就记录任务名”:`tasks_map[task_obj.name] = future`,这样完全不依赖 API,就不会再遇到 ObjectNotFound。 - 如果你必须用 API 拿 TaskRun.name,请确保运行和查询使用同一个 Prefect API 配置(同一 profile/workspace,且 flow 非离线运行)。 参考 - Futures/状态:prefect.futures - Orchestration 客户端:orchestration client - State 与异常提取:prefect.states 如果你贴一下当前运行环境(是否连接 API、profile 名称)和你提交任务时的代码块,我可以把方案 A 的映射整合到你的
my_flow
里,帮你把返回结构也统一好。
r
@Marvin def task(___fn_=None, **_kwargs_): """task""" # 统一设定部分参数默认值 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): wrapped = prefect_task(**kwargs)(fn) return wrapped # 支持 @task 和 @task(...) return _decorate if __fn is None else _decorate(__fn) 这是我自定义的task装饰器,我希望默认情况下task_run_name和函数名保持一致,怎么写代码
m
thought for 27.9 seconds
可以,把默认的 task_run_name 设为“函数名”就行。你自定义装饰器里在包装函数处读取 `fn.__name__`,用
setdefault
赋给
name
task_run_name
即可(用户手动传入时仍能覆盖)。 给你一版修正且支持 @task 和 @task(...) 的实现(顺便修了你变量名的小问题):
Copy code
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
换成可调用:
Copy code
kwargs.setdefault(
      "task_run_name",
      lambda ctx, params: fn.__name__  # 这里可以根据 params 组装更丰富的名字
  )
r
@Marvin task_runs = await client.read_task_runs( _task_run_filter_=TaskRunFilter(_flow_run_id_={"any_": [flow_run.id]})) api获取的taskrun里怎么获取task的函数名
m
thought for 133.7 seconds
结论:通过 Orchestration API 拿到的
TaskRun
里没有“函数名”这个字段。能用的是 `task_key`,它的格式是“函数名-代码哈希”。因此你可以从
task_run.task_key
里解析出函数名部分。 实用做法 - 直接解析 `task_key`:
Copy code
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]
- 放到你的查询代码里用:
Copy code
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
@Marvin def get_failed_tasks(_local_vars_): """ 获取失败的任务名称 """ print("True Order:", list(local_vars.keys())) failed_var = "" failed_msg = "" for var_name, var_val in local_vars.items(): print(f"check_val:{var_name}, type:{type(var_val)}") # 检查 PrefectFuture if isinstance(var_val, PrefectFuture): print(f"find_PrefectFuture:{var_name}") st = var_val.state if st.is_failed(): print(f"find_failed:{var_name}") failed_var = var_name failed_msg = get_err_msg(st) break # 检查 PrefectFutureList elif isinstance(var_val, PrefectFutureList): print(f"find_PrefectFutureList:{var_name}") if len(var_val) == 0: continue for index, fut in enumerate(var_val): st = fut.state print(f"check_type:{type(fut)}, st:{st}") if st.is_failed(): print(f"find_failed:{var_name}") failed_var = var_name failed_msg = get_err_msg(st) break if failed_var != "": break return failed_var, failed_msg 这是我通过遍历task返回值检查任务失败状态的函数,但检查PrefectFutureList这个分支里,st = fut.state的值是Pending(),导致无法检查,怎么解决
m
thought for 30.0 seconds
出现 Pending 的原因是你在检查时任务还没完成。PrefectFuture 的
state
只有在任务进入终态(COMPLETED/FAILED/CRASHED/CANCELLED)后才会反映最终结果;未完成时就是 PENDING/RUNNING。 解决思路有三种,任选其一即可: 方案 A:先等待所有 futures 完成,再检查状态(推荐)
Copy code
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 完成 - 对单个元素:
Copy code
val = fut.result(raise_on_failure=False)  # 等到完成,不会抛异常;失败时返回异常对象
st = fut.state  # 此时为终态
- 对列表:
Copy code
for fut in var_val:
    fut.result(raise_on_failure=False)
    if fut.state.is_failed():
        ...
方案 C:用 as_completed 按完成顺序扫描,遇到第一个失败就返回
Copy code
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 容易误伤其他变量,建议提交时维护一个专门的映射,比如:
Copy code
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.states
r
@Marvin 每次flow执行完后调用def get_failed_tasks(local_vars): """ 获取失败的任务名称 """ failed_var = "" failed_msg = "" 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)]) _, _ = wait(all_futs, timeout=60) if all_futs else (set(), set()) # done, not_done for var_name, var_val in local_vars.items(): # 检查 PrefectFuture if isinstance(var_val, PrefectFuture): st = var_val.state if st.is_failed(): failed_var = var_name failed_msg = get_err_msg(st) break # 检查 PrefectFutureList elif isinstance(var_val, PrefectFutureList): if len(var_val) == 0: continue for index, fut in enumerate(var_val): st = fut.state if st.is_failed(): failed_var = var_name failed_msg = get_err_msg(st) break if failed_var != "": break 是否会导致进程占用的线程数暴增,pgsql压力增大
m
thought for 218.9 seconds
不会。 基于 Prefect 3.x 的实现细节,下面是对你这段逻辑对线程数和数据库负载影响的结论与建议: 结论 - 线程数不会暴增 -
prefect.futures.wait()
使用事件驱动的等待机制:进程内会有一个共享的异步事件循环/订阅器,而不是为每个 future 新建一个线程。对成千上万个 futures 调用一次
wait()
也不会产生成比例的新线程。 - 不会造成持续的 PostgreSQL 压力 -
wait()
对远程运行的任务会做一次“初次读取”检查(每个 future 至少一次 API 读取任务状态),之后通过事件流(websocket)等待状态变化,不做高频轮询。 - 访问
future.state
时,如果状态未缓存且未终结,会触发一次按需的 API 读取;状态一旦终结,会在本地缓存,之后反复访问不会再次打 API。 - 真正会放大压力的反例 - 频繁循环访问
.state
并 sleep 的做法会造成高频 API 读取(类似轮询),应避免:
Copy code
# 反例:会造成高频 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,避免无关变量混入。 更高效/更稳妥的写法示例 - 一次性等待 + 批量取状态(避免不必要的结果读取):
Copy code
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 的压力峰值,还是担心潜在风险?