如何配置Apache Beam将S3作为artifacts_dir存储目录?
解决Apache Beam PortableRunner下artifacts_dir配置S3不生效的问题
针对你使用Apache Beam 2.41.0 + Flink 1.14.5 + PortableRunner时,配置artifacts_dir到S3不生效、扩展服务启动失败的问题,按以下步骤排查和解决:
1. 确保S3文件系统依赖完整
Beam的S3存储依赖Hadoop的s3a客户端,必须确保作业和JobServer(扩展服务)都包含对应依赖:
- Java环境:作业构建时引入
org.apache.beam:beam-sdks-java-io-hadoop-file-system:2.41.0,同时JobServer的classpath需要包含hadoop-aws、aws-java-sdk-bundle(版本需和Beam兼容,Beam 2.41.0适配Hadoop 3.2.x)。 - Python环境:安装时指定
apache-beam[aws]==2.41.0,确保AWS相关依赖被加载。
2. 正确配置JobServer(扩展服务)参数
PortableRunner的工件暂存由JobServer管理,必须在JobServer启动时指定artifacts_dir,而非仅在作业提交时配置。启动命令示例:
java -jar beam-runners-flink-1.14-job-server-2.41.0.jar \ --flink-master=localhost:8081 \ --artifacts_dir=s3a://my-bucket/beam-artifact-staging \ --job-port=8099 \ --artifact-port=8098
同时需配置S3访问凭证:
- 通过环境变量:
AWS_ACCESS_KEY_ID和AWS_SECRET_ACCESS_KEY - 或在Hadoop
core-site.xml中添加:
并将<property> <name>fs.s3a.access.key</name> <value>你的AWS Access Key</value> </property> <property> <name>fs.s3a.secret.key</name> <value>你的AWS Secret Key</value> </property>core-site.xml放入JobServer的classpath目录。
3. 作业提交时同步配置参数
提交作业时需再次指定artifacts_dir,并传递文件系统配置:
Java作业示例
java -jar your-pipeline.jar \ --runner=FlinkPortableRunner \ --jobEndpoint=localhost:8099 \ --artifacts_dir=s3a://my-bucket/beam-artifact-staging \ --filesystemConfiguration=fs.s3a.access.key=xxx,fs.s3a.secret.key=xxx
Python作业示例
python your_pipeline.py \ --runner=FlinkPortableRunner \ --job_endpoint=localhost:8099 \ --artifacts_dir=s3a://my-bucket/beam-artifact-staging \ --filesystem_configuration=fs.s3a.access.key=xxx,fs.s3a.secret.key=xxx
4. 排查扩展服务启动失败问题
若JobServer添加artifacts_dir后无法启动,优先检查:
- 依赖缺失:使用官方构建的全依赖JobServer jar,或自行构建包含所有S3相关依赖的fat jar。
- S3配置错误:查看JobServer启动日志,确认是否有凭证无效、bucket不存在、权限不足等错误。
- 参数格式:确保
artifacts_dir路径为合法的s3a://bucket/path格式,bucket需提前创建。
5. 验证配置生效
作业启动后,可通过以下方式验证:
- 查看Flink UI的作业日志,确认无本地
/tmp/beam-artifact-staging的相关日志。 - 检查指定的S3 bucket,确认生成了
beam-artifact-staging目录及相关工件文件。
内容的提问来源于stack exchange,提问作者Lydian
相关产品推荐
相关产品推荐

