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

如何用Django+Celery持续监控Redis列表?求合理实现方案

Django+Celery持续监控Redis列表的优化方案

方案1:阻塞式轮询+信号处理(解决CTRL+C无法终止问题)

替换原有的while True+sleep轮询逻辑,用Redis的BLPOP阻塞命令减少无效请求,同时注册信号处理函数响应终止信号,让任务能被CTRL+C正常停止。

修改后的任务代码:

from celery import shared_task
import redis
import signal

should_stop = False

def handle_termination(signum, frame):
    global should_stop
    should_stop = True

@shared_task(bind=True)
def watcher(self):
    # 注册SIGINT信号处理,捕获CTRL+C
    signal.signal(signal.SIGINT, handle_termination)
    
    client = redis.Redis('localhost', port=6379, decode_responses=True)
    queue = 'queue'
    lock_key = 'watcher_key'

    try:
        if client.get(lock_key):
            print('监控任务已在运行')
            return

        # 给锁设置过期时间,避免异常退出后锁残留
        client.set(lock_key, 'true', ex=30)
        print(f'启动Redis列表 {queue} 监控...')

        while not should_stop:
            # 阻塞等待列表元素,超时1秒(保证能及时响应终止信号)
            # 若需保留列表元素,可改用brpoplpush将元素临时转移后再放回
            result = client.blpop([queue], timeout=1)
            if result:
                _, item = result
                print(f'列表 {queue} 新增元素: {item}')
                # 如需保留元素,取消下面注释:
                # client.lpush(queue, item)

    except Exception as e:
        print(f'Redis错误: {e}')
    finally:
        client.delete(lock_key)
        print('监控任务已停止')

方案2:复用Celery Worker,合并Beat启动(无需额外终端)

不用单独开终端启动django-celery-beat,直接用命令让Worker和Beat在同一个终端运行,将监控任务改为周期性任务:

from celery import shared_task
from celery.schedules import crontab
from celery.task import periodic_task
import redis

# 每秒执行一次监控
@periodic_task(run_every=crontab(minute='*', second='*/1'))
def watcher():
    client = redis.Redis('localhost', port=6379, decode_responses=True)
    queue = 'queue'
    lock_key = 'watcher_key'

    # 用Redis锁防止多Worker实例重复执行任务
    with client.lock(lock_key, timeout=10):
        try:
            current_items = client.lrange(queue, 0, -1)
            print(f'列表 {queue} 当前元素: {current_items}')
        except Exception as e:
            print(f'Redis错误: {e}')

启动命令:

celery -A your_django_project worker -B --loglevel=info

这样只需要两个终端:一个跑Django服务器,一个跑合并了Beat的Celery Worker。

方案3:Redis发布订阅(Pub/Sub)模式(实时监控变化)

如果只需要在列表有更新时触发监控,而非定时轮询,可以用Redis的Pub/Sub机制:在往列表添加元素的代码里同时发布消息,监控任务订阅对应频道,实时响应变化。

生产者代码(添加元素时)

client = redis.Redis('localhost', port=6379, decode_responses=True)
# 添加元素到列表
client.lpush('queue', 'new_item')
# 发布更新消息到频道
client.publish('queue_update_channel', 'item_added')

监控任务代码

from celery import shared_task
import redis

@shared_task(bind=True)
def watcher(self):
    client = redis.Redis('localhost', port=6379, decode_responses=True)
    pubsub = client.pubsub()
    pubsub.subscribe('queue_update_channel')

    lock_key = 'watcher_key'
    if client.get(lock_key):
        print('监控任务已在运行')
        return
    client.set(lock_key, 'true', ex=30)

    try:
        print('开始订阅列表更新频道...')
        for message in pubsub.listen():
            if message['type'] == 'message':
                # 收到更新消息后获取最新列表内容
                current_items = client.lrange('queue', 0, -1)
                print(f'列表 {queue} 更新后元素: {current_items}')
            # 响应任务终止信号
            if self.request.called_directly and getattr(self, 'should_stop', False):
                break
    except Exception as e:
        print(f'Redis错误: {e}')
    finally:
        pubsub.unsubscribe()
        client.delete(lock_key)
        print('监控任务已停止')

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 03:33:30