You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Prefect 1.0如何为SQLiteQuery等内置任务添加日志记录

为Prefect 1.0预置任务添加日志的实现方法

在Prefect 1.0版本中,@task装饰器定义的自定义任务可通过prefect.context.get("logger")获取上下文日志实例输出内容,但prefect.tasks.database.sqlite.SQLiteQuery这类官方预置任务类默认没有内置日志逻辑,直接实例化不会产生执行日志。以下是两种符合1.0版本规范的实现方案:


方案1:继承预置任务类,自定义带日志的子类(推荐复用场景使用)

直接继承原预置任务类,重写run方法插入日志逻辑,原有参数、功能完全兼容,可全局复用:

from prefect import task, Flow
import prefect
from time import sleep
from prefect.tasks.database.sqlite import SQLiteQuery

# 扩展原SQLiteQuery类,添加日志逻辑
class LoggedSQLiteQuery(SQLiteQuery):
    def run(self, *args, **kwargs):
        logger = prefect.context.get("logger")
        logger.info(f"开始执行SQLite查询 | 数据库路径: {self.db} | 执行SQL: {self.query}")
        try:
            result = super().run(*args, **kwargs)
            logger.info(f"SQLite查询执行完成,返回结果: {result}")
            return result
        except Exception as e:
            logger.error(f"SQLite查询执行失败,错误信息: {str(e)}")
            raise e

@task()
def some_task():
    logger = prefect.context.get("logger")
    logger.info("Let's sleep a second!")
    sleep(1)

# 实例化带日志能力的查询任务
version_check = LoggedSQLiteQuery(
    db="sqlite.db",
    query="Select sqlite_version()",
)


with Flow("a flow") as flow:
    some_task()
    # 注意:原示例仅打印任务实例、未加()调用,任务不会被Flow调度执行
    query_result = version_check()


if __name__ == "__main__":
    flow.run()

方案2:轻量包装调用(适合单次使用场景)

不需要新建子类,直接在Flow内用@task包装一层调用逻辑即可,代码改动量最小:

from prefect import task, Flow
import prefect
from time import sleep
from prefect.tasks.database.sqlite import SQLiteQuery


@task()
def some_task():
    logger = prefect.context.get("logger")
    logger.info("Let's sleep a second!")
    sleep(1)


version_check = SQLiteQuery(
    db="sqlite.db",
    query="Select sqlite_version()",
)

# 包装任务,内部调用预置任务的run方法并加日志
@task
def exec_version_query():
    logger = prefect.context.get("logger")
    logger.info("开始执行SQLite版本查询")
    res = version_check.run()
    logger.info(f"查询完成,当前SQLite版本: {res[0][0]}")
    return res


with Flow("a flow") as flow:
    some_task()
    db_version = exec_version_query()


if __name__ == "__main__":
    flow.run()

补充说明

  • 两种方式获取的logger和自定义@task中的日志实例完全一致,会自动携带任务名、Flow运行ID等Prefect内置日志字段,日志格式、级别配置全局统一
  • 所有Prefect 1.0预置任务(如ShellTask、各类云服务请求任务等)都可以用上述两种方式添加自定义日志,逻辑通用
  • 不要直接修改Prefect库源码添加日志,后续版本升级会丢失自定义修改,继承/包装的方式兼容性更好

内容的提问来源于stack exchange,提问作者ivegotaquestion

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 20:24:29