Apache Beam Python SDK跨语言读取MQTT报错排查求助
解决方案
核心原因
Beam自动下载的默认beam-sdks-java-expansion-service-app-2.50.0.jar仅包含核心模块,MqttIO属于beam-sdks-java-io-mqtt扩展IO模块,未被包含在默认jar中,因此扩展服务找不到该类。
解决办法
方法1:手动启动包含MqttIO依赖的扩展服务
- 确保本地安装Maven
- 执行以下命令启动扩展服务(指定
io-mqttprofile引入依赖):
mvn exec:java -Dexec.mainClass=org.apache.beam.sdk.expansion.service.ExpansionService -Dexec.args="8091" -Pio-mqtt -Dbeam.version=2.50.0
- 修改Python代码,指定扩展服务地址:
- 在创建
JavaExternalTransform时添加expansion_service参数:read_java_transform = JavaExternalTransform( 'org.apache.beam.sdk.io.mqtt.MqttIO', expansion_service='localhost:8091' ).read().withConnectionConfiguration(...) - 或者在
PipelineOptions中设置:beam_options = PipelineOptions(beam_args, expansion_service_url='localhost:8091')
- 在创建
方法2:在Python代码中直接指定依赖
修改JavaExternalTransform的初始化参数,添加classpath引入MqttIO依赖:
read_java_transform = JavaExternalTransform( 'org.apache.beam.sdk.io.mqtt.MqttIO', classpath=['org.apache.beam:beam-sdks-java-io-mqtt:2.50.0'] ).read().withConnectionConfiguration(...)
Beam会自动下载对应依赖并加载到扩展服务中。
额外代码修正
注意:withConnectionConfiguration需要接收MqttIO.ConnectionConfiguration对象,原代码用list(mqtt_config)传递参数方式不正确,应改为:
# 替换原MqttConfig和mqtt_config的定义 mqtt_config = { 'serverUri': server_uri, 'topic': topic, 'clientId': 'random123', 'username': username, 'password': password } read_java_transform = JavaExternalTransform( 'org.apache.beam.sdk.io.mqtt.MqttIO', classpath=['org.apache.beam:beam-sdks-java-io-mqtt:2.50.0'] ).read().withConnectionConfiguration(**mqtt_config)
内容的提问来源于stack exchange,提问作者Daniel Stevenson
相关产品推荐
相关产品推荐

