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

Django Channels异步消息无法实时发送问题及优化方案咨询

问题:Django Channels 异步接收数据后无法实时发送的优化方案咨询

问题背景

我有一个通过httpx异步接收数据的业务流程,期望数据一被接收就通过Django Channels发送,但Channels却会在所有数据接收完成后才统一发送。我使用的是内存版channel layer,无法安装Redis服务,尝试了多线程、run_in_executor等方案均未解决问题。

复现代码

class AsyncSend(AsyncJsonWebsocketConsumer):

    async def connect(self):
        # Join group
        self.group_name = f"test"

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

        await self.accept()

        #Replicate non blocking call
        async for result in self.non_block():
            await self.send_task(result)

    async def non_block(self):
        for integer in range(4):
            await asyncio.sleep(1)
            curTime = datetime.now().strftime("%H:%M:%S")
            print(f"{curTime} - {integer}")
            yield integer

    async def send_task(self, result):
        curTime = datetime.now().strftime("%H:%M:%S")
        print(f"{curTime} - Send task called")
        await self.channel_layer.group_send(
            self.group_name, {
                'type': "transmit",
                'message': result,
            }
        )
        print("Finished sending")

    async def transmit(self, event):
        print(f"Received event to transmit...")

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

问题输出结果

10:23:21 - 0
10:23:21 - Send task called
Finished sending
10:23:22 - 1
10:23:22 - Send task called
Finished sending
10:23:23 - 2
10:23:23 - Send task called
Finished sending
10:23:24 - 3
10:23:24 - Send task called
Finished sending
Received event to transmit...
Received event to transmit...
Received event to transmit...
Received event to transmit...

尝试过的无效方案

  • 多线程运行Consumer的方案
  • 使用run_in_executor和ThreadPoolExecutor的实现
    以上方案均未解决问题,执行顺序仍是串行的。

当前可行方案

发现问题在于Channels会等待函数执行完成后才处理队列中的其他任务,于是将调用外部流程的代码从connect方法移至receive_json方法,并通过asyncio.create_task让receive_json立即返回,实现了预期的实时发送效果:

修改后的关键代码

async def receive_json(self, text_data):
    asyncio.create_task(self.submit())

async def submit(self):
    #Replicate non blocking call
    async for result in self.non_block():
        await self.send_task(result)

优化后输出结果

15:16:56 - 0
15:16:56 - Send task called
Finished sending
Received 0 to transmit...
15:16:57 - 1
15:16:57 - Send task called
Finished sending
Received 1 to transmit...
15:16:58 - 2
15:16:58 - Send task called
Finished sending
Received 2 to transmit...
15:16:59 - 3
15:16:59 - Send task called
Finished sending
Received 3 to transmit...

提问

想咨询是否有更优的实现方式来达成这一需求?


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 20:30:49