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

Azure ServiceBus Python SDK消费消息时触发TypeError错误求助

Azure ServiceBus Python SDK消费消息时TypeError问题解决方案

问题描述

使用Azure ServiceBus Python异步SDK消费队列消息,在处理部分消息后触发以下错误:

TypeError: Received message 257:None is not bytes

完整错误栈:

socket.send() raised exception.
Unexpected error occurred (TypeError('Received message 257:None is not bytes')). Handler shutting down.
Traceback (most recent call last):
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/site-packages/azure/servicebus/aio/_base_handler_async.py", line 269, in _do_retryable_operation
    return await operation(**kwargs)
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/site-packages/azure/servicebus/aio/_servicebus_receiver_async.py", line 410, in _receive
    receiving = await amqp_receive_client.do_work_async()
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/site-packages/azure/servicebus/_pyamqp/aio/_client_async.py", line 357, in do_work_async
    return await self._client_run_async(**kwargs)
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/site-packages/azure/servicebus/_pyamqp/aio/_client_async.py", line 753, in _client_run_async
    await self._connection.listen(wait=self._socket_timeout, **kwargs)
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/site-packages/azure/servicebus/_pyamqp/aio/_connection_async.py", line 798, in listen
    if await self._read_frame(wait=wait, **kwargs):
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/site-packages/azure/servicebus/_pyamqp/aio/_connection_async.py", line 268, in _read_frame
    new_frame = await self._transport.receive_frame(timeout=timeout, **kwargs)
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/site-packages/azure/servicebus/_pyamqp/aio/_transport_async.py", line 76, in receive_frame
    header, channel, payload = await asyncio.wait_for(
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/asyncio/tasks.py", line 445, in wait_for
    return fut.result()
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/asyncio/futures.py", line 201, in result
    raise self._exception.with_traceback(self._exception_tb)
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/asyncio/tasks.py", line 232, in __step
    result = coro.send(None)
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/site-packages/azure/servicebus/_pyamqp/aio/_transport_async.py", line 133, in read
    await self._read(payload_size, buffer=payload)
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/site-packages/azure/servicebus/_pyamqp/aio/_transport_async.py", line 532, in _read
    data = await self.sock.receive_bytes()
  File "/home/trusted-service-user/cluster-env/clonedenv/lib/python3.10/site-packages/aiohttp/client_ws.py", line 281, in receive_bytes
    raise TypeError(f"Received message {msg.type}:{msg.data!r} is not bytes")
TypeError: Received message 257:None is not bytes
exception:Handler failed: Received message 257:None is not bytes.

使用的消费代码:

from azure.servicebus import TransportType
from azure.servicebus.aio import ServiceBusClient
servicebus_client = ServiceBusClient.from_connection_string(connection_string,
                                                            transport_type=TransportType.AmqpOverWebsocket,
                                                            max_lock_duration=180, logging_enable=True)
async with servicebus_client:
    receiver = servicebus_client.get_queue_receiver(queue_name)
    async with receiver:
       messages = await receiver.receive_messages(max_message_count=5000, max_wait_time=60)

环境信息:Python 3.10,azure-servicebus版本7.11.2

解决方案

1. 升级SDK版本

该错误是SDK内部AMQP over Websocket传输层的已知问题,已在7.12.0及以上版本修复。执行升级命令:

pip install --upgrade azure-servicebus>=7.12.0

2. 调整接收参数(临时 workaround)

若暂时无法升级SDK,可降低max_message_count值(例如从5000调整为1000),减少单次接收的消息数量,降低触发异常的概率。

3. 添加异常捕获与重连逻辑

在消费代码中增加异常捕获,当出现该错误时自动重试,避免进程崩溃:

from azure.servicebus import TransportType, ServiceBusError
from azure.servicebus.aio import ServiceBusClient
import asyncio

async def consume_messages(connection_string, queue_name):
    while True:
        try:
            servicebus_client = ServiceBusClient.from_connection_string(
                connection_string,
                transport_type=TransportType.AmqpOverWebsocket,
                max_lock_duration=180, 
                logging_enable=True
            )
            async with servicebus_client:
                receiver = servicebus_client.get_queue_receiver(queue_name)
                async with receiver:
                    messages = await receiver.receive_messages(max_message_count=1000, max_wait_time=60)
                    # 处理消息逻辑
                    for msg in messages:
                        await receiver.complete_message(msg)
        except (TypeError, ServiceBusError) as e:
            print(f"消费出错,将重试: {str(e)}")
            await asyncio.sleep(5)  # 等待5秒后重试

# 调用消费函数
# asyncio.run(consume_messages("<你的连接字符串>", "<队列名>"))

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 13:35:13