Django: как подружить Celery с WebSocket

Короткий ответ: в Django Celery и WebSocket связывают через Channels и channel layer. Celery выполняет фоновую задачу, отправляет событие в группу Channels, а WebSocket consumer получает это событие и передает его браузеру. Celery не должен напрямую держать WebSocket-соединение с пользователем

Схема:

Browser -> WebSocket -> Django Channels consumer
Celery task -> channel layer -> group_send -> consumer -> Browser

Что нужно установить

Обычно нужны:

pip install channels channels-redis celery redis

Redis часто используют и как брокер Celery, и как backend для channel layer. Настройки зависят от проекта, но идея одна: Django, Celery и Channels должны видеть общий Redis

Настройка channel layer

INSTALLED_APPS = [
    "channels",
    # ваши приложения
]

ASGI_APPLICATION = "project.asgi.application"

CHANNEL_LAYERS = {
    "default": {
        "BACKEND": "channels_redis.core.RedisChannelLayer",
        "CONFIG": {
            "hosts": [("127.0.0.1", 6379)],
        },
    },
}

WebSocket consumer

import json
from channels.generic.websocket import AsyncWebsocketConsumer

class ProgressConsumer(AsyncWebsocketConsumer):
    async def connect(self):
        self.user_id = self.scope["url_route"]["kwargs"]["user_id"]
        self.group_name = f"user_{self.user_id}_progress"

        await self.channel_layer.group_add(self.group_name, self.channel_name)
        await self.accept()

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

    async def progress_message(self, event):
        await self.send(text_data=json.dumps({
            "type": "progress",
            "payload": event["payload"],
        }))

Метод progress_message вызывается, когда в группу отправляют событие с type: "progress.message"

Celery task отправляет событие

from asgiref.sync import async_to_sync
from channels.layers import get_channel_layer
from celery import shared_task

@shared_task
def process_report(user_id):
    channel_layer = get_channel_layer()
    group_name = f"user_{user_id}_progress"

    for percent in [10, 30, 60, 100]:
        async_to_sync(channel_layer.group_send)(
            group_name,
            {
                "type": "progress.message",
                "payload": {"percent": percent},
            },
        )

    return {"status": "done"}

Celery не знает про конкретный WebSocket. Он отправляет событие в группу, а Channels доставляет его подключенным клиентам

Где запускать задачу

Задачу можно запустить из обычного Django view:

from django.http import JsonResponse
from .tasks import process_report

def start_report(request):
    process_report.delay(request.user.id)
    return JsonResponse({"status": "started"})

Пользователь нажимает кнопку, HTTP-запрос запускает Celery task, а прогресс приходит уже через WebSocket. Это нормальное разделение: HTTP стартует действие, Celery выполняет тяжелую работу, WebSocket показывает live-статус

Что хранить в базе

Если задача важная, не полагайтесь только на WebSocket-события. Сохраняйте итоговый статус в базе:

pending -> processing -> done -> failed

Тогда пользователь, который обновил страницу или потерял соединение, сможет снова открыть интерфейс и увидеть актуальное состояние. WebSocket в этой схеме ускоряет доставку обновлений, но база остается источником правды

Клиент в браузере

const socket = new WebSocket("ws://localhost:8000/ws/progress/42/");

socket.addEventListener("message", (event) => {
  const message = JSON.parse(event.data);

  if (message.type === "progress") {
    console.log("progress:", message.payload.percent);
  }
});

Частые ошибки

Первая ошибка — пытаться отправлять WebSocket-сообщение прямо из Celery в браузер. У Celery task нет браузерного соединения

Вторая ошибка — забыть group_add в consumer. Тогда Celery отправляет событие, но никто его не получает

Третья ошибка — неправильно указать type события. В Channels progress.message попадет в метод progress_message

Четвертая ошибка — использовать разные Redis или разные настройки channel layer у Django и Celery

Пятая ошибка — не проверять права. Пользователь должен получать только свои события, а не чужой прогресс

Шестая ошибка — считать WebSocket надежным хранилищем статуса. Если событие потерялось из-за обрыва соединения, пользователь должен иметь возможность получить состояние через обычный HTTP-запрос

Самопроверка

Откройте WebSocket-клиент для пользователя 42, запустите Celery task process_report.delay(42) и проверьте, что в браузер приходят проценты. Затем запустите задачу для другого user_id и убедитесь, что первый клиент не получает чужие сообщения

Что почитать дальше по WebSocket

Если нужен общий маршрут по теме, откройте рубрику WebSocket. Для соседних задач пригодятся эти разборы:

Оцените статью
0 0 голоса
Рейтинг статьи
Подписаться
Уведомить о
guest

0 комментариев
Старые
Новые Популярные
Межтекстовые Отзывы
Посмотреть все комментарии
0
Оставьте комментарий! Напишите, что думаете по поводу статьи.x