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
相关产品推荐
相关产品推荐

