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

如何用Django Channels与Celery实现定时消息推送?避免重复启动任务

实现首个用户触发的定时推送任务方案

核心思路

要实现仅在有用户连接WebSocket时启动定时推送任务,无在线用户时自动停止,且避免重复启动任务,相比始终运行的定时任务,这种方案能有效节省服务器资源,是更优的选择。我们可以通过缓存标记任务状态、动态启停Celery循环任务来实现。

具体实现代码

修改后的 consumers.py

import json
from django.core.cache import cache
from channels.generic.websocket import AsyncWebsocketConsumer
from app_monitor.tasks import monitor_loop


class MonitorConsumer(AsyncWebsocketConsumer):
    GROUP_NAME = 'monitor'
    # 缓存键,用于标记定时任务是否在运行
    TASK_RUNNING_FLAG = 'monitor_task_active'

    async def connect(self):
        await self.channel_layer.group_add(self.GROUP_NAME, self.channel_name)
        await self.accept()

        # 检查缓存,无运行标记则启动循环任务
        if not cache.get(self.TASK_RUNNING_FLAG):
            # 启动任务,间隔10秒推送一次(可根据业务调整)
            monitor_loop.delay(self.GROUP_NAME, 10)
            cache.set(self.TASK_RUNNING_FLAG, True, timeout=None)

    async def disconnect(self, close_code):
        await self.channel_layer.group_discard(self.GROUP_NAME, self.channel_name)

        # 查询当前群组的在线连接数
        active_channels = await self.channel_layer.group_channels(self.GROUP_NAME)
        # 无在线用户时,清除任务运行标记,任务会自动终止
        if not active_channels:
            cache.delete(self.TASK_RUNNING_FLAG)

    async def receive(self, text_data):
        text_data_json = json.loads(text_data)
        message = text_data_json['message']

    async def send_data(self, event):
        await self.send(text_data=json.dumps({"data": event['data']}))

修改后的 tasks.py

from dbaOps.celery import app
from channels.layers import get_channel_layer
from asgiref.sync import async_to_sync
from django.db import connections
from django.core.cache import cache
from app_monitor.utils.sqlarea import MONITOR_SQL


@app.task(bind=True)
def monitor_loop(self, group, interval):
    # 先检查缓存标记,无标记则终止任务循环
    if not cache.get('monitor_task_active'):
        return

    try:
        with connections['tools'].cursor() as cur:
            cur.execute(MONITOR_SQL)
            data = cur.fetchall()
        channel_layer = get_channel_layer()
        async_to_sync(channel_layer.group_send)(
            group,
            {'type': 'send.data', 'data': data}
        )
    except Exception as e:
        print(e)

    # 递归调用自身,实现定时循环推送
    self.apply_async(args=[group, interval], countdown=interval)

方案优势

  • 资源高效:无用户连接时任务自动停止,避免空跑查询和推送操作
  • 避免重复:缓存标记确保同一时刻只有一个定时任务在运行,不会因多用户连接启动多个任务
  • 动态适配:用户连接时自动启动任务,最后一个用户断开后自动停止,无需人工干预

注意事项

  • 确保缓存后端(如Redis)稳定运行,任务启停依赖缓存标记
  • 根据业务需求调整任务间隔,避免过于频繁查询数据库造成压力
  • 若任务执行时间超过设定间隔,可适当调大间隔值,避免任务堆积

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 04:43:25