使用Python Proton连接Eclipse Hono AMQP Adaptor发送遥测数据报错求助
问题描述
我目前尝试通过AMQP Adaptor向Hono沙箱发送遥测消息,虽然我参考了Hono北向桥接示例中的部分代码(理论上南向桥接也可使用该方案),但似乎在SASL认证环节遇到了问题。
我的第一版代码如下:
from __future__ import print_function, unicode_literals from proton import Message from proton.handlers import MessagingHandler from proton.reactor import Container tenantId = 'xxxx' deviceId = 'yyyyy' devicePassword = 'my-secret-password' class AmqpMessageSender(MessagingHandler): def __init__(self, server, address): super(AmqpMessageSender, self).__init__() self.server = server self.address = address def on_start(self, event): conn = event.container.connect( self.server, sasl_enabled=True, allowed_mechs="PLAIN", allow_insecure_mechs=True, user=f'{deviceId}@{tenantId}', password=devicePassword ) event.container.create_sender(conn, self.address) def on_sendable(self, event): msg = Message( address=f'{self.address}/{deviceId}', content_type='application/json', body={"temp": 5, "transport": "amqp"} ) event.sender.send(self.msg) event.sender.close() def on_connection_error(self, event): print("Connection Error") def on_link_error(self, event): print("Link Error") def on_transport_error(self, event): print("Transport Error") Container(AmqpMessageSender(f'amqp://hono.eclipseprojects.io:5671', f'telemetry/{tenantId}')).run()
运行代码时出现传输错误,上下文提示信息为:
'Expected SASL protocol header: no protocol header found (connection aborted)'
我也尝试使用5672端口,返回链路错误;使用实际为北向桥接端口的15672时,意外没有触发SASL错误,但返回了预期的“未授权”错误(因为设备不允许通过北向桥接连接)。
更新内容
再次感谢您抽出时间解答。由于评论区字数有限,我在此重新贴出用于模拟设备的代码:
from __future__ import print_function, unicode_literals from proton import Message from proton.handlers import MessagingHandler from proton.reactor import Container tenantId = 'xxx' deviceId = 'yyy' devicePassword = 'my-secret-password' class AmqpMessageSender(MessagingHandler): def __init__(self, server): super(AmqpMessageSender, self).__init__() self.server = server def on_start(self, event): print("In start") conn = event.container.connect( self.server, sasl_enabled=True, allowed_mechs="PLAIN", allow_insecure_mechs=True, user=f'{deviceId}@{tenantId}', password=devicePassword ) print("connection established") event.container.create_sender(context=conn, target=None) print("sender created") def on_sendable(self, event): print("In Msg send") event.sender.send(Message( address=f'telemetry', properties={ 'to': 'telemetry', 'content-type': 'application/json' }, content_type='application/json', body={"temp": 5, "transport": "amqp"} )) event.sender.close() event.connection.close() print("Sender & connection closed") def on_connection_error(self, event): print("Connection Error") def on_link_error(self, event): print("Link Error") def on_transport_error(self, event): print("Transport Error") Container(AmqpMessageSender(f'amqp://hono.eclipseprojects.io:5672')).run()
我没有使用Java客户端模拟服务端,而是同样采用Python快速入门示例的代码。我还编写了客户端类可按快速入门示例发起HTTP调用,服务端类可响应调用并打印消息,因此我认为下方的服务端实现没有问题:
from __future__ import print_function, unicode_literals import threading import time from proton.handlers import MessagingHandler from proton.reactor import Container amqpNetworkIp = "hono.eclipseprojects.io" tenantId = 'xxx' class AmqpReceiver(MessagingHandler): def __init__(self, server, address, name): super(AmqpReceiver, self).__init__() self.server = server self.address = address self._name = name def on_start(self, event): conn = event.container.connect(self.server, user="consumer@HONO", password="verysecret") event.container.create_receiver(conn, self.address) def on_connection_error(self, event): print("Connection Error") def on_link_error(self, event): print("Link Error") def on_message(self, event): print(self._name) print("Got a message:") print(event.message.body) class CentralServer: def listen_telemetry(self, name): uri = f'amqp://{amqpNetworkIp}:15672' address = f'telemetry/{tenantId}' self.container = Container(AmqpReceiver(uri, address, name)) print("Starting (northbound) AMQP Connection...") self.thread = threading.Thread(target=lambda: self.container.run(), daemon=True) self.thread.start() time.sleep(2) def stop(self): # Stop container print("Stopping (northbound) AMQP Connection...") self.container.stop() self.thread.join(timeout=5) CentralServer().listen_telemetry('cs1')
我又尝试了一天仍未定位到问题,非常希望您能指出我遗漏的部分:)
此致
Armin
解决方案
你代码里存在三处核心错误,修正后即可正常运行:
- 协议与端口不匹配:Hono沙箱的5671端口是TLS加密的南向AMQP端口,你使用
amqp://前缀访问会跳过TLS握手,服务端返回的TLS报文被识别为无效SASL头,直接触发你遇到的传输错误。如果要使用5671端口,需要将前缀改为amqps://;如果使用非加密的5672端口,保留amqp://即可。 - 发送端目标地址配置错误:Hono南向AMQP规范要求,设备上报遥测时必须指定sender的target地址为
telemetry,你更新后的代码将target设为None,会直接返回链路错误。 - 代码语法错误:第一版代码的
on_sendable方法中你定义了局部变量msg,但发送时调用了不存在的self.msg,会直接触发运行时异常。
以下是修正后的可运行设备端代码,使用非加密的5672端口即可:
from __future__ import print_function, unicode_literals from proton import Message from proton.handlers import MessagingHandler from proton.reactor import Container tenantId = 'xxx' # 替换为你实际的租户ID deviceId = 'yyy' # 替换为你实际的设备ID devicePassword = 'my-secret-password' # 替换为你实际的设备密码 class AmqpMessageSender(MessagingHandler): def __init__(self, server): super(AmqpMessageSender, self).__init__() self.server = server def on_start(self, event): print("In start") conn = event.container.connect( self.server, sasl_enabled=True, allowed_mechs="PLAIN", allow_insecure_mechs=True, user=f'{deviceId}@{tenantId}', password=devicePassword ) print("connection established") # 修正:指定target为telemetry event.container.create_sender(context=conn, target="telemetry") print("sender created") def on_sendable(self, event): print("In Msg send") event.sender.send(Message( content_type='application/json', body={"temp": 5, "transport": "amqp"} )) event.sender.close() event.connection.close() print("Sender & connection closed") def on_connection_error(self, event): print("Connection Error") # 可打印远程返回的具体错误原因 if event.connection.remote_condition: print(f"错误原因:{event.connection.remote_condition.description}") def on_link_error(self, event): print("Link Error") if event.link.remote_condition: print(f"错误原因:{event.link.remote_condition.description}") def on_transport_error(self, event): print("Transport Error") if event.transport.remote_condition: print(f"错误原因:{event.transport.remote_condition.description}") # 非加密端口用amqp://前缀 Container(AmqpMessageSender(f'amqp://hono.eclipseprojects.io:5672')).run()
验证步骤
- 确保你在Hono沙箱中注册的租户、设备、设备凭证配置正确,密码和代码中的
devicePassword完全一致。 - 先启动你编写的北向消费端代码,确认无报错、正常运行。
- 运行修正后的设备端代码,即可在消费端看到上报的遥测数据。
内容的提问来源于stack exchange,提问作者Armin Gruber
相关产品推荐
相关产品推荐

