在AWS Kinesis Zeppelin中使用PyFlink初始化环境遇AttributeError求助
问题排查:AttributeError: 'EnvironmentSettings' has no attribute 'in_streaming_mode'
核心原因
- AWS Kinesis Data Analytics(KDA)当前绑定的PyFlink版本大概率低于1.12,而
EnvironmentSettings.in_streaming_mode()是Flink 1.12及以上版本才新增的静态API,低版本环境中不存在该方法。 - 你参考的官方示例可能基于较高版本的Flink,但KDA的运行环境版本滞后,导致API不兼容。
解决方法
1. 最简流环境初始化(推荐)
低版本PyFlink中,StreamExecutionEnvironment.get_execution_environment()默认就是流处理模式,无需额外指定:
from pyflink.datastream import StreamExecutionEnvironment env = StreamExecutionEnvironment.get_execution_environment()
2. 自定义环境配置的写法
如果需要配置并行度、状态后端等参数,直接通过StreamExecutionEnvironment的方法设置:
env = StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(3) # 其他配置比如状态后端设置 # env.set_state_backend(...)
3. 兼容低版本的EnvironmentSettings用法
若必须通过EnvironmentSettings定义环境模式,低版本需要先实例化对象再链式调用方法:
from pyflink.datastream import StreamExecutionEnvironment from pyflink.common import EnvironmentSettings env_settings = EnvironmentSettings.new_instance().in_streaming_mode().build() env = StreamExecutionEnvironment.get_execution_environment(env_settings)
有效参考文档
- 先在KDA控制台确认当前使用的PyFlink版本,再查阅对应版本的Flink官方PyFlink文档,重点看流环境初始化章节。
- 优先参考AWS官方的Kinesis Data Analytics PyFlink开发文档,其中会明确列出支持的API和版本匹配规则,避免套用通用Flink最新示例。
内容的提问来源于stack exchange,提问作者hi im Bacon
相关产品推荐
相关产品推荐

