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

使用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()

验证步骤

  1. 确保你在Hono沙箱中注册的租户、设备、设备凭证配置正确,密码和代码中的devicePassword完全一致。
  2. 先启动你编写的北向消费端代码,确认无报错、正常运行。
  3. 运行修正后的设备端代码,即可在消费端看到上报的遥测数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 18:27:02