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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 06:35:19