如何用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ô
相关产品推荐
相关产品推荐

