Apache Beam MQTTIO连接不同Broker对应主题时表现不一致问题咨询
问题根因判断
你遇到的低频率消息仅能接收第一条的问题,大概率是新Broker的空闲超时配置与MQTTIO默认的心跳配置不匹配导致的,可通过以下步骤验证并解决问题:
第一步:核对新Broker的
keepalive_timeout配置
多数MQTT Broker默认空闲超时为60秒,恰好匹配你提到的「每分钟1条」的低频率场景。你之前使用的Mosquitto默认空闲超时通常为10分钟,所以5分钟间隔的消息不会触发断连,更换Broker后阈值变小就会出现该问题:当客户端超过超时时间未发送心跳包或业务消息,Broker会主动断开连接,若MQTTIO默认未开启自动重连,后续消息就无法被接收。第二步:补充MQTTIO连接的心跳与重连配置
在ConnectionConfiguration中新增以下参数即可解决:ConnectionConfiguration config = ConnectionConfiguration.create(newBrokerUri, clientId) // 心跳间隔设置为小于Broker空闲超时的一半,比如Broker超时60秒就设为25秒 .withKeepAliveInterval(25) // 开启自动重连,断连后自动恢复订阅 .withAutomaticReconnect(true) // 关闭Clean Session,重连后可接收断连期间的离线消息(需配合QoS≥1使用) .withCleanSession(false);第三步:核对消息QoS配置
确保你向新Broker发布消息的QoS等级≥1,若使用QoS 0,断连期间的消息会直接丢失,不会被补发。补充说明:你之前配置的
withMaxNumRecords、withMaxReadTime仅限制单次读取的最大条数或最长时间,和长连接保活无关,无法解决断连导致的消息丢失问题;--streaming=true参数仅控制Beam的运行模式,也不影响MQTT连接层的保活逻辑,因此调整这两个配置无效是正常现象。
内容的提问来源于stack exchange,提问作者Viswadeep V
相关产品推荐
相关产品推荐

