FastAPI环境下Prefect 2.0.x MinIO Storage Block加载/保存失败问题
FastAPI环境中使用Prefect Storage Block触发JSONDecodeError问题
错误栈
File "/usr/lib/python3.8/json/__init__.py", line 357, in loads return _default_decoder.decode(s) File "/usr/lib/python3.8/json/decoder.py", line 337, in decode obj, end = self.raw_decode(s, idx=_w(s, 0).end()) File "/usr/lib/python3.8/json/decoder.py", line 355, in raw_decode raise JSONDecodeError("Expecting value", s, err.value) from None json.decoder.JSONDecodeError: Expecting value: line 1 column 1 (char 0) Expecting value: line 1 column 1 (char 0)
环境配置
- Prefect Server:
prefecthq/prefect:2.20-python3.8 - Prefect Worker:
prefecthq/prefect:2.20-python3.8(已安装prefect-aws库) - MinIO:
minio/minio:RELEASE.2024-08-17T01-24-54Z
相关代码
service.py
from prefect import get_client from prefect_aws.s3 import S3Bucket import importlib, traceback from prefect.deployments.deployments import Deployment class PrefectProvider: def __init__(self): self.prefect_client = get_client() async def register_flow(self, flow_name: str, schedule): try: minio_storage = await S3Bucket.load("minio-s3") flow_fun = importlib.import_module(f"flows.{flow_name}") deployment = await Deployment.build_from_flow( flow=getattr(flow_fun, flow_name), flow_name=flow_name, name=f"{flow_name}-deployment", schedules=[schedule], infra_overrides={"env": {"PREFECT_LOGGING_LEVEL": "DEBUG"}}, storage=minio_storage, work_queue_name="minio-worker" ) deployment_id = await deployment.apply() return deployment_id except Exception as e: traceback.print_exc() raise e
main.py
from fastapi import FastAPI, Depends from service import PrefectProvider import uvicorn from prefect.client.schemas.schedules import CronSchedule from module import prefect_provider_di app = FastAPI() @app.get("/") async def root(provider: PrefectProvider = Depends(prefect_provider_di)): await provider.register_flow("my_test_flow", CronSchedule(cron="* * * * *")) if __name__ == "__main__": uvicorn.run(app, host="0.0.0.0", port=8000)
module.py
from service import PrefectProvider def prefect_provider_di(): return PrefectProvider()
flows/my_test_flow.py
from prefect import flow @flow def my_test_flow(): print("Hello guys") if __name__ == "__main__": my_test_flow()
问题现象
- 调用
S3Bucket.load()或新建Block并执行save()时触发上述JSONDecodeError - 若调用
load时不添加await,会抛出coroutine 'Block.load' was never awaited错误 - VSCode调试模式下代码可正常运行
- 重启PC后问题自行解决,但未明确具体原因
内容的提问来源于stack exchange,提问作者Galih Anggara
相关产品推荐
相关产品推荐

