文档: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)