文档:https://docs.celeryq.dev/en/stable/index.html
pip install celery
pip install celery-once 防止任务重复,目前用不上
协程支持
pip install gevent
基础知识教程
开发时的数据库注意
celery不支持异步方法,仅支持同步,因此要时刻注意维护两个数据库操作版本,各自的连接相互不要影响,包括调用的文件,不要把对方的数据库导入(本来celery就是单独的进程,它自己需要一个单独的连接池)
并发模式选择
- 进程:CPU 密集型任务(如数据处理、图片压缩、视频转码)多进程独立运行,避免 GIL 限制
- 线程:有少量 I/O 等待,轻微 CPU 操作 适合轻量的任务,但受 Python GIL 限制,不适合高并发或密集型工作
- 协程:I/O 密集型任务(如网络请求、API调用、数据库操作),高并发、资源占用低,适合处理大量等待网络响应的任务,不适合 CPU 密集任务
示例
from celery import Celery
from celery.schedules import crontab
import pytz # 正确使用时区
import os
import dotenv
from datetime import datetime
dotenv.load_dotenv()
REDIS_HOST = os.getenv("REDIS_HOST")
REDIS_PASSWORD = os.getenv("REDIS_PASSWORD")
# 定义一个时区常量
SHANGHAI_TZ = pytz.timezone('Asia/Shanghai')
# 初始化 Celery 应用,使用 Redis 作为 broker
app = Celery("celery_app", broker=f"redis://:{REDIS_PASSWORD}@{REDIS_HOST}:6379/0")
# 设置全局时区为北京时间
app.conf.timezone = "Asia/Shanghai"
# 生命周期
@app.on_after_configure.connect
def periodic_tasks(sender: Celery, **kwargs):
print("Celery消费者启动")
# 动态配置周期任务,通过sender.add_periodic_task动态添加
#这里的30也可以用也可以:crontab(minute='*/10')每十分钟执行一次
sender.add_periodic_task(30, test.s('每30秒执行一次'), name='30秒任务',expires=10)
# 每周一早上7点30分执行
sender.add_periodic_task(
crontab(hour=7, minute=30, day_of_week=1),
test.s('Happy Mondays!'),
name='每周任务'
)
# 设置补发任务时间只检查最近30秒,防止周期任务补发任务
app.conf.beat_cron_starting_deadline = 30
# 静态配置周期任务,在启动时就配置无法中途修改
app.conf.beat_schedule = {
'20秒任务': {
'task': 'add_task', # 静态配置任务必须为任务名或完整模块路径,例如:crowdpulse.celery_ser.celery_app.add_task
'schedule': 20,
'args': (16, 16)
},
}
# 申明任务, 申明后可以通过任务函数名.delay()立即调用任务
# 最好将完整模块路径作为name,例如crowdpulse.celery_ser.celery_app.test
@app.task(name='test')
def test(arg):
print(arg)
# 申明一个多参数任务,并显示指定任务名
@app.task(name='add_task')
def add_task(x, y):
z = x + y
print(f"20秒任务执行结果:{z}")
if __name__ == "__main__":
# 立即执行任务
test.delay("立即执行任务")
# 10秒后执行
test.apply_async(args=["10秒后执行任务"], countdown=10)
# 今天下午13点47分执行
# 转换为上海时间
target_time = SHANGHAI_TZ.localize(datetime(2025, 5, 7, 13, 47))
# 开始定时任务
test.apply_async(args=["今天下午13点47分执行任务"], eta=target_time)
优先级
redis
# 使用redis需设置
app.conf.broker_transport_options = {
'queue_order_strategy': 'priority',
}
# 这里的name最好写完整的模块路径
@app.task(name='test')
def test(data):
print(f"执行任务:{data}")
# 测试
if __name__ == "__main__":
# 立即执行任务
test.delay("立即执行任务")
# 立即执行任务,优先级9
test.apply_async(args=["立即执行,优先级9"], priority=9)
test.apply_async(args=["立即执行,优先级5"], priority=5)
test.apply_async(args=["立即执行,优先级0"], priority=0)
redis对优先级支持不好,作用小,要支持优先级推荐RabbitMQ
RabbitMQ(存在eta长时间关闭问题)
pip install celery[pyamqp] 可能需安装依赖
手动安装文档:
https://www.rabbitmq.com/docs/install-debian
容器安装
sudo docker pull rabbitmq:4.1-management
运行容器
sudo docker run -d --name my-rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:4.1-management
如果需要设置心跳超时时间毫秒 在docker run后面加上-e RABBITMQ_SERVER_ADDITIONAL_ERL_ARGS="-rabbit consumer_timeout 7200000"
然后访问主机IP+15672端口访问管理界面,默认管理账号密码:guest
app = Celery(
'my_project',
broker='pyamqp://guest@localhost//', # RabbitMQ 地址
backend='rpc://' # 可选,用于获取任务结果
)
注意事项
在RabbitMQ中队列一旦创建,然后Python改了新的队列配置,需要将RabbitMQ中的队列删除,不然可能不会更新配置,点击队列名称进去后删除即可
清空还未被领取的任务
点击进入队列后,点击Purge按钮会清空所有未被执行的任务
队列页面的列表说明:
Ready 等待被取走的任务
Unacked 已被取走但未执行完的任务
Total 总任务数
celery-once 防止任务重复
注意它控制的是相同参数在同一个任务中不重复执行,如果你的任务实际中就需要同样的参数执行多次,则不宜使用这个工具
# 在app实列中需要
app = Celery("celery_app", broker=f"redis://:{REDIS_PASSWORD}@{REDIS_HOST}:6379/0")
app.conf.ONCE = {
'backend': 'celery_once.backends.Redis',
'settings': {
'url': f'redis://:{REDIS_PASSWORD}@{REDIS_HOST}:6379/0',
'default_timeout': 600 #如果600秒没有释放锁,则其它任务能重试它
}
}
# 在实际的任务中需要
from celery_once import QueueOnce
@app.task(base=QueueOnce,name='publish_twitter')
def publish_twitter(username:str,content:str):
# 假设这里是一个任务
pass
路由:为不同任务指定不同消费者
app.conf.update(# 设置默认队列为 ai
task_default_queue='ai',
task_routes={
'tasks.ai_tasks.*': {'queue': 'ai'}, #交给协程
'tasks.browser_tasks.*': {'queue': 'browser'}, # 交给线程
})
消费者启动命令
AI任务专用worker:协程池高并发
celery -A crowdpulse.celery_app worker -n ai_worker@%h -Q ai -P gevent -c 100
浏览器任务worker:多进程隔离执行
celery -A crowdpulse.celery_app worker -n browser_worker@%h -Q browser -P prefork -c 8
消费者服务启动
启动消费者命令
celery -A crowdpulse.celery_ser.celery_app worker --loglevel=info
启动周期任务命令
celery -A crowdpulse.celery_ser.celery_app beat
同时启动周期任务和消费者
多并发时生产环境不建议
celery -A crowdpulse.celery_ser.celery_app worker -B
消费者高并发处理
线程启动
20个线程启动
celery -A crowdpulse.celery_ser.celery_app worker --pool=threads --concurrency=20
协程启动
100 个协程
celery -A crowdpulse.celery_ser.celery_app worker -P gevent -c 100
多进程协程
多进程,它们自己会向redis拉取任务
启动第一个进程,100协程
celery -A crowdpulse.celery_ser.celery_app worker -P gevent -c 100 -n worker1@%h
启动第二个进程,100协程
celery -A crowdpulse.celery_ser.celery_app worker -P gevent -c 100 -n worker2@%h
系统服务启动消费者进程
注意:因为服务即进程,所以要启动多个进程服务,就要使用多个系统启动服务来执行进程,同时名字和配置名称也记得区分开
/etc/systemd/system/celery.service
[Unit]
Description=celery
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User=ubuntu
Group=ubuntu
WorkingDirectory=/home/ubuntu/crowdpulse
Environment="PATH=/home/ubuntu/crowdpulse/.venv/bin"
# 等待 Redis 启动并能连接
ExecStartPre=/bin/bash -c 'for i in {1..30}; do nc -z 192.168.1.11 6379 && break || sleep 1; done'
ExecStart=/home/ubuntu/crowdpulse/.venv/bin/celery -A crowdpulse.celery_ser.celery_app worker -B
Restart=always
RestartSec=10
[Install]
WantedBy=multi-user.target
直接在python代码中使用subprocess运行消费者进程
不推荐这样做
import subprocess
import os
# 设置虚拟环境路径和工作目录
venv_path = "/home/ubuntu/crowdpulse/.venv/bin"
work_dir = "/home/ubuntu/crowdpulse"
# 构造 Celery 启动命令
celery_command = [
os.path.join(venv_path, "celery"),
"-A",
"crowdpulse.celery_ser.celery_app",
"worker",
"-B"
]
# 启动 Celery worker
subprocess.Popen(
celery_command,
cwd=work_dir, # 设置工作目录
env={**os.environ, "PATH": venv_path}, # 设置环境变量
# start_new_session=True, python主程序结束后,消费者依然存在
# stdout=subprocess.PIPE, # 标准输出到管道,除非主动读取,不然不会打印
# stderr=subprocess.PIPE # 标准错误到管道,除非主动读取,不然不会打印
)
如果要传对象只使用可序列化对象
社区推荐方法是:只传数据的主键id,到任务中再执行一次查询
例如数据库读取的orm参数这些都是不可用的,可以先自己做一层转换:
# 转换为支持点语法访问的序列化对象
persona_obj = {
"id": persona.id,
"name": persona.name,
"bio": persona.bio,
"traits": persona.traits,
}
task_fun.delay(persona_obj) # 传入序列化对象
@app.task(name='task_fun')
def task_fun(persona_dict):
# 将字典转换为支持点语法访问的对象
class DotDict(dict):
"""支持点语法访问的字典类"""
def __getattr__(self, attr):
return self.get(attr)
__setattr__ = dict.__setitem__
__delattr__ = dict.__delitem__
# 转换为支持点语法访问的对象
persona = DotDict(persona_dict)
# 接下来可以使用persona.bio对象了