AWS托管Apache Flink执行权限拒绝问题求助
问题:AWS托管Apache Flink运行PyFlink应用时权限拒绝错误
环境配置
- Apache Flink 1.18
- Kinesis connector 1.18
- Python
apache-flink == 1.18 - 通过pom.xml收集
flink-connector-kinesis 4.2.0-1.18的Jar依赖,将uber-jar保存至pyflink/lib/ - Python 3.10/3.8均出现相同错误,本地运行正常,AWS中启动任务时触发权限拒绝
极简复现代码
from pyflink.datastream.stream_execution_environment import StreamExecutionEnvironment def main(): env = StreamExecutionEnvironment.get_execution_environment() try: env.from_collection(['a', 'b', 'c', 'd']).map(lambda x: x).map(print) env.execute('executethis') except Exception as ex: print('EXCEPTION!!!! ', ex) if __name__ == '__main__': main()
错误堆栈信息
java.lang.RuntimeException: Failed to create stage bundle factory! at org.apache.flink.streaming.api.runners.python.beam.BeamPythonFunctionRunner.createStageBundleFactory(BeamPythonFunctionRunner.java:656) at org.apache.flink.streaming.api.runners.python.beam.BeamPythonFunctionRunner.open(BeamPythonFunctionRunner.java:281) at org.apache.flink.streaming.api.operators.python.process.AbstractExternalPythonFunctionOperator.open(AbstractExternalPythonFunctionOperator.java:57) at org.apache.flink.streaming.api.operators.python.process.AbstractExternalDataStreamPythonFunctionOperator.open(AbstractExternalDataStreamPythonFunctionOperator.java:85) at org.apache.flink.streaming.api.operators.python.process.AbstractExternalOneInputPythonFunctionOperator.open(AbstractExternalOneInputPythonFunctionOperator.java:117) at org.apache.flink.streaming.api.operators.python.process.ExternalPythonProcessOperator.open(ExternalPythonProcessOperator.java:64) at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.initializeStateAndOpenOperators(RegularOperatorChain.java:107) at org.apache.flink.streaming.runtime.tasks.StreamTask.restoreGates(StreamTask.java:753) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.call(StreamTaskActionExecutor.java:100) at org.apache.flink.streaming.runtime.tasks.StreamTask.restoreInternal(StreamTask.java:728) at org.apache.flink.streaming.runtime.tasks.StreamTask.restore(StreamTask.java:693) at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:955) at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:924) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:748) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:564) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.UncheckedExecutionException: java.io.IOException: Cannot run program "/tmp/python-dist-ea8c2605-8024-4255-a9d0-f9432a304ee7/python-files/py_site_packages38/py_site_packages38/pyflink/bin/pyflink-udf-runner.sh": error=13, Permission denied at org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.LocalCache$LocalLoadingCache.getUnchecked(LocalCache.java:5022) at org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory$SimpleStageBundleFactory.<init>(DefaultJobBundleFactory.java:498) at org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory$SimpleStageBundleFactory.<init>(DefaultJobBundleFactory.java:482) at org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory.forStage(DefaultJobBundleFactory.java:342) at org.apache.flink.streaming.api.runners.python.beam.BeamPythonFunctionRunner.createStageBundleFactory(BeamPythonFunctionRunner.java:654) 15 more **Caused by: java.io.IOException: Cannot run program "/tmp/python-dist-ea8c2605-8024-4255-a9d0-f9432a304ee7/python-files/py_site_packages38/py_site_packages38/pyflink/bin/pyflink-udf-runner.sh": error=13, Permission denied **at java.base/java.lang.ProcessBuilder.start(ProcessBuilder.java:1128) at java.base/java.lang.ProcessBuilder.start(ProcessBuilder.java:1071) at org.apache.beam.runners.fnexecution.environment.ProcessManager.startProcess(ProcessManager.java:147) at org.apache.beam.runners.fnexecution.environment.ProcessManager.startProcess(ProcessManager.java:122) at org.apache.beam.runners.fnexecution.environment.ProcessEnvironmentFactory.createEnvironment(ProcessEnvironmentFactory.java:104) at org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory$1.load(DefaultJobBundleFactory.java:284) at org.apache.beam.runners.fnexecution.control.DefaultJobBundleFactory$1.load(DefaultJobBundleFactory.java:240) at org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.LocalCache$LoadingValueReference.loadFuture(LocalCache.java:3571) at org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.LocalCache$Segment.loadSync(LocalCache.java:2313) at org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.LocalCache$Segment.lockedGetOrLoad(LocalCache.java:2190) at org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.LocalCache$Segment.get(LocalCache.java:2080) at org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.LocalCache.get(LocalCache.java:4012) at org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.LocalCache.getOrLoad(LocalCache.java:4035) at org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.LocalCache$LocalLoadingCache.get(LocalCache.java:5013) at org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.LocalCache$LocalLoadingCache.getUnchecked(LocalCache.java:5020) 19 more Suppressed: java.lang.NullPointerException: Process for id does not exist: 7-1 at org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull(Preconditions.java:921) at org.apache.beam.runners.fnexecution.environment.ProcessManager.stopProcess(ProcessManager.java:172) at org.apache.beam.runners.fnexecution.environment.ProcessEnvironmentFactory.createEnvironment(ProcessEnvironmentFactory.java:124) 29 more Caused by: java.io.IOException: error=13, Permission denied at java.base/java.lang.ProcessImpl.forkAndExec(Native Method) at java.base/java.lang.ProcessImpl.<init>(ProcessImpl.java:340) at java.base/java.lang.ProcessImpl.start(ProcessImpl.java:271) at java.base/java.lang.ProcessBuilder.start(ProcessBuilder.java:1107)
已尝试的解决方案
- 为IAM策略添加所有Kinesis Analytics权限
- 为代码所在S3前缀添加所有S3权限
- 为输入输出Kinesis流添加所有权限
- 为IAM策略添加所有Kinesis Analytics V2权限
- 按照官方文档添加CloudWatch相关权限
- 将本地包权限设置为777后再打包成应用ZIP
- 确保Jar版本与AWS托管Flink 1.18匹配,限制Python依赖为1.18版本
- 降级Python至3.8,确保所有Python库对应Flink 1.18版本
- 删除所有冗余代码及Kinesis相关逻辑
- 添加try/except块捕获异常
- 将所有依赖库打包到uber-jar中,并在配置中引用
解决建议
1. 确保打包时保留脚本可执行权限
问题核心是pyflink-udf-runner.sh无执行权限,即使本地设置了777,打包ZIP时可能丢失权限:
- 执行
chmod +x pyflink/bin/*.sh为所有shell脚本添加可执行权限 - 使用
zip -r -y app.zip .命令打包(-y参数保留文件权限和符号链接)
2. 配置Flink Python临时目录
AWS托管Flink的/tmp目录可能存在权限限制,可指定其他有权限的临时路径:
在Flink配置中添加:
python.tmp.dir=/mnt/tmp/pyflink
/mnt目录在AWS托管节点通常有足够读写执行权限。
3. 替换lambda映射避免触发UDF逻辑
map(lambda x: x)会被Flink识别为Python UDF,触发外部进程。可改用显式类型声明的内置算子:
from pyflink.common.typeinfo import Types env.from_collection(['a', 'b', 'c', 'd'], Types.STRING()) \ .map(lambda x: x, output_type=Types.STRING()) \ .print()
或直接简化逻辑,避免触发外部Python进程。
4. 使用AWS官方提供的Python依赖
避免自行打包pyflink依赖,直接使用AWS托管Flink兼容的官方Python运行环境,确保依赖包权限符合平台要求。
内容的提问来源于stack exchange,提问作者elrorris
相关产品推荐
相关产品推荐

