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

如何实现Django与Angular14的Socket连接?解决Celery任务状态获取问题

Django + Angular 14 实现Socket连接与Celery任务状态高效获取

一、解决Django与Angular 14的Socket连接问题(针对ngx-socket-io失败场景)

Django原生不支持WebSocket,需借助channels实现后端WebSocket服务,再配合ngx-socket-io完成前端对接,步骤如下:

后端Django配置

  1. 安装依赖
pip install channels channels-redis django-cors-headers
  1. 修改settings.py配置
    • 将channels和corsheaders加入INSTALLED_APPS
    • 设置ASGI应用:ASGI_APPLICATION = 'your_project_name.asgi.application'
    • 配置通道层(以Redis为例):
      CHANNEL_LAYERS = {
          "default": {
              "BACKEND": "channels_redis.core.RedisChannelLayer",
              "CONFIG": {
                  "hosts": [("127.0.0.1", 6379)],
              },
          },
      }
      
    • 配置CORS允许WebSocket跨域:
      CORS_ALLOW_ALL_ORIGINS = True  # 生产环境建议指定具体域名
      CORS_ALLOW_CREDENTIALS = True
      
  2. 编写WebSocket消费者
    创建consumers.py,定义处理任务状态推送的逻辑:
import json
from channels.generic.websocket import AsyncWebsocketConsumer

class TaskStatusConsumer(AsyncWebsocketConsumer):
    async def connect(self):
        # 加入全局任务更新分组,也可按用户/任务组自定义分组
        self.group_name = "task_updates"
        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 task_status_update(self, event):
        await self.send(text_data=json.dumps({
            "task_id": event["task_id"],
            "status": event["status"],
            "result": event.get("result"),
            "error": event.get("error")
        }))
  1. 配置WebSocket路由
    在项目根路由文件(如urls.py)中添加:
from django.urls import path
from . import consumers

websocket_urlpatterns = [
    path("ws/tasks/", consumers.TaskStatusConsumer.as_asgi()),
]
  1. 启动ASGI服务器
    替换Django默认的runserver,用daphne启动:
daphne your_project_name.asgi:application

前端Angular 14调整

  1. 确保ngx-socket-io安装正确
npm install ngx-socket-io
  1. 配置Socket连接
    在app.module.ts中配置后端WebSocket地址:
import { NgModule } from '@angular/core';
import { SocketModule, SocketIoConfig } from 'ngx-socket-io';

// 注意使用ws/wss协议,地址对应后端WebSocket路由
const config: SocketIoConfig = { 
    url: 'ws://localhost:8000/ws/tasks/', 
    options: { withCredentials: true } 
};

@NgModule({
    imports: [SocketModule.forRoot(config)],
    ...
})
export class AppModule { }
  1. 创建Socket服务监听状态更新
// src/app/services/socket.service.ts
import { Injectable } from '@angular/core';
import { Socket } from 'ngx-socket-io';
import { Observable } from 'rxjs';

@Injectable({ providedIn: 'root' })
export class SocketService {
    // 订阅后端推送的任务状态更新
    taskStatus$: Observable<any> = this.socket.fromEvent('task_status_update');

    constructor(private socket: Socket) {}
}
  1. 在组件中使用服务
    在需要展示任务状态的组件中注入SocketService,订阅状态更新并刷新UI:
import { Component, OnInit } from '@angular/core';
import { SocketService } from './services/socket.service';

@Component({ ... })
export class TaskStatusComponent implements OnInit {
    taskStatuses: Map<string, any> = new Map();

    constructor(private socketService: SocketService) {}

    ngOnInit(): void {
        this.socketService.taskStatus$.subscribe(update => {
            this.taskStatuses.set(update.task_id, update);
        });
    }
}

二、Celery任务状态高效获取方案(替代轮询)

利用Celery信号监听任务状态变化,通过WebSocket主动推送到前端,无需轮询:

  1. 配置Celery信号
    在Celery配置文件(如celery.py)中添加信号监听:
from celery import Celery
from celery.signals import after_task_publish, task_success, task_failure
from channels.layers import get_channel_layer
from asgiref.sync import async_to_sync

app = Celery('your_project', broker='redis://localhost:6379/0')

channel_layer = get_channel_layer()

# 任务发布时发送状态
@after_task_publish.connect
def notify_task_sent(sender=None, headers=None, **kwargs):
    task_id = headers['id']
    async_to_sync(channel_layer.group_send)(
        "task_updates",
        {
            "type": "task_status_update",
            "task_id": task_id,
            "status": "PENDING"
        }
    )

# 任务成功时发送状态
@task_success.connect
def notify_task_success(sender=None, result=None, **kwargs):
    task_id = sender.request.id
    async_to_sync(channel_layer.group_send)(
        "task_updates",
        {
            "type": "task_status_update",
            "task_id": task_id,
            "status": "SUCCESS",
            "result": result
        }
    )

# 任务失败时发送状态
@task_failure.connect
def notify_task_failure(sender=None, exception=None, **kwargs):
    task_id = sender.request.id
    async_to_sync(channel_layer.group_send)(
        "task_updates",
        {
            "type": "task_status_update",
            "task_id": task_id,
            "status": "FAILURE",
            "error": str(exception)
        }
    )
  1. 额外优化
  • 若需按用户隔离任务状态,可在消费者连接时获取用户身份,加入用户专属分组,信号发送时指定对应分组
  • 对于长耗时任务,可在任务执行过程中主动发送中间状态更新(通过调用channel_layer发送消息)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 15:45:39