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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 08:26:33