使用FlinkKafkaConsumer的PyFlink集群任务报‘No module named google’错误
PyFlink集群运行报错:ModuleNotFoundError: No module named 'google'
我正在开发一个PyFlink任务,通过FlinkKafkaConsumer连接器从Kafka Topic读取数据。该任务在本地运行正常,但提交至Flink集群时,持续出现与google模块相关的报错,具体错误信息如下:
File "/tmp/pyflink/7fc48e92-9dc7-468b-9c34-cddc89128621/efe4c9e6-a4e9-41be-a7da-cf7fe8be3165/flink_local.py", line 35, in main stream.map(write_to_file, output_type=Types.STRING()) File "/usr/lib/flink/opt/python/pyflink.zip/pyflink/datastream/data_stream.py", line 291, in map File "/usr/lib/flink/opt/python/pyflink.zip/pyflink/datastream/data_stream.py", line 557, in process File "<frozen zipimport>", line 259, in load_module File "/usr/lib/flink/opt/python/pyflink.zip/pyflink/fn_execution/flink_fn_execution_pb2.py", line 23, in <module> ModuleNotFoundError: No module named 'google' org.apache.flink.client.program.ProgramAbortException: java.lang.RuntimeException: Python process exits with code: 1
已尝试的排查步骤:
- 确保虚拟环境中已安装所需依赖库
- 确认PyFlink任务运行在正确的虚拟环境中
- 检查Flink配置,确保虚拟环境设置正确
- 验证flink_local.py脚本中的导入语句准确一致
- 本地测试脚本运行正常
尽管尝试了以上方法,仍无法解决该错误。任务代码如下:
from pyflink.datastream import StreamExecutionEnvironment from pyflink.datastream.connectors import FlinkKafkaConsumer from pyflink.common.serialization import SimpleStringSchema from pyflink.common.typeinfo import Types from pyflink.common import RestartStrategies def write_to_file(value): with open('flink-output.txt', 'a') as f: f.write(value + '\n') return value def main(): env = StreamExecutionEnvironment.get_execution_environment() env.set_restart_strategy(RestartStrategies.fixed_delay_restart(3, 1000)) properties = { "bootstrap.servers": "localhost:9092", "group.id": "test-group", } schema = SimpleStringSchema() kafka_consumer = FlinkKafkaConsumer("flink-read", schema, properties) kafka_consumer.set_start_from_earliest() stream = env.add_source(kafka_consumer) stream.map(write_to_file, output_type=Types.STRING()) env.execute("Kafka to File") if __name__ == '__main__': main()
内容的提问来源于stack exchange,提问作者sam1064max
相关产品推荐
相关产品推荐

