celery任务队列


文档: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 密集任务

示例

python
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

python
# 使用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

python
app = Celery(
    'my_project',
    broker='pyamqp://guest@localhost//',  # RabbitMQ 地址
    backend='rpc://'  # 可选,用于获取任务结果
)

注意事项
在RabbitMQ中队列一旦创建,然后Python改了新的队列配置,需要将RabbitMQ中的队列删除,不然可能不会更新配置,点击队列名称进去后删除即可

清空还未被领取的任务
点击进入队列后,点击Purge按钮会清空所有未被执行的任务
队列页面的列表说明:
Ready 等待被取走的任务
Unacked 已被取走但未执行完的任务
Total 总任务数

celery-once 防止任务重复

注意它控制的是相同参数在同一个任务中不重复执行,如果你的任务实际中就需要同样的参数执行多次,则不宜使用这个工具

python
# 在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

路由:为不同任务指定不同消费者

python
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

ini
[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运行消费者进程

不推荐这样做

python
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参数这些都是不可用的,可以先自己做一层转换:

python
# 转换为支持点语法访问的序列化对象
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对象了