Prefect 任务编排+队列


文档:https://docs.prefect.io/v3/get-started

pip install prefect

服务器
提供web ui监控面板
存储和跟踪任务状态
分发任务给消费者进程

prefect worker消费者进程
实际执行任务的进程,会不断向prefect server拉取任务

任务流flow
能被prefect调用的最小单元,它能够编排任务

任务task
所有任务都要被放入flow才能被执行

步骤:
prefect server start 启动服务器
prefect config view 查看配置

消费者程序,很多时候不需要用到,因为更多时候只用来编排
prefect work-pool create default -t process 创建work pool
prefect worker start -p default 启动work

python crowdpulse/celery_ser/prefect_demo.py执行

python
from prefect import flow, task
from datetime import datetime
from prefect.task_worker import serve

@task(log_prints=True)
def hello(name: str):
    print(f"{name},你好")


@flow(name="hello_flow")
def hello_flow(name: str):
    hello(f"运行在工作流中的:{name}")


# 
if __name__ == "__main__":
    # # 立即运行任务流,会阻塞程序,一个任务流是一个进程
    # hello_flow.serve(
    #     name="hello_flow",
    #     cron="* * * * *",
    #     parameters={"name": "李坤朋"}
    # )
    hello.delay("李坤朋") # 开始一个任务到队列
    hello.delay("桦一") # 开始第二任务到队列
    hello.serve() # 启动一个任务工作器,并执行等待的任务
    # # 第二种方法
    # hello.map("李坤朋2", deferred=True) # background 3 task runs - i.e. zip(A, B)
    # hello.map("桦一2", deferred=True) # background 3 task runs - i.e. zip(A, B)
    # serve(hello)