Apache Beam Java SDK访问Cloud Storage时出现401无效凭证错误
Apache Beam访问GCS存储桶出现401认证错误问题
问题场景
尝试从Cloud Storage存储桶读取CSV文件并存入PCollection,使用拥有roles/storage.admin角色的服务账号JSON密钥做认证,但Apache Beam管道运行时持续报401「无效凭证」错误。
Beam Pipeline配置代码
DataflowPipelineOptions dfOptions = PipelineOptionsFactory.as(DataflowPipelineOptions.class); dfOptions.setProject("project_name"); dfOptions.setStagingLocation("bucket_name"); dfOptions.setGcpCredential(GoogleCredentials.fromStream( new FileInputStream(PATH_TO_JSON_KEY))); dfOptions.setTempLocation("gs://bucket_name/folder_name"); dfOptions.setServiceAccount("serivce_acount_name"); Pipeline myPipe= Pipeline.create(dfOptions); PCollection<ReadableFile> readFile= myPipe.apply( FileIO.match().filepattern("gs://bucket_name/file_name.csv")).apply(FileIO.readMatches());
错误信息
- 本地运行时错误:
Caused by: java.io.IOException: Error trying to get gs://bucket_name/object_name.csv: {"code":401,"errors":[{"domain":"global","location":"Authorization","locationType":"header","message":"Invalid Credentials","reason":"authError"}],"message":"Invalid Credentials"}
- 添加
DataflowRunner后,针对暂存桶的错误:
401 Unauthorized GET https://storage.googleapis.com/storage/v1/b/bucket_name { "code" : 401, ...same as above... }
验证结果
- 相同凭证通过GCS Java客户端库访问存储桶正常:
StorageOptions options = StorageOptions.newBuilder() .setProjectId(PROJECT_ID) .setCredentials(GoogleCredentials.fromStream( new FileInputStream(PATH_TO_JSON_KEY))).build(); Storage storage = options.getService(); Blob blob = storage.get(BUCKET_NAME, OBJECT_NAME); ReadChannel r = blob.reader();
- 相同服务账号密钥通过
gsutil下载文件无问题,仅Beam场景报错。
依赖版本
- Apache Beam 2.24
- google-cloud-storage 2.11.3
- google-cloud-dataflow-java-sdk-all 2.5.0
- google-api-client 1.35.1
解决建议
1. 修复配置中的明显错误
- StagingLocation路径格式错误:
setStagingLocation必须传入完整的gs://bucket_name/path格式路径,当前仅传bucket_name会导致路径解析失败,进而触发认证错误。 - ServiceAccount拼写错误:代码中
serivce_acount_name拼写错误,需修正为正确的服务账号邮箱格式(如xxx@project-id.iam.gserviceaccount.com)。
2. 避免凭证配置冲突
Beam中同时设置gcpCredential和serviceAccount可能导致冲突,建议二选一:
方式一:仅使用JSON密钥文件(本地运行优先)
移除setServiceAccount配置,确保凭证加载时添加必要的权限范围:
DataflowPipelineOptions dfOptions = PipelineOptionsFactory.as(DataflowPipelineOptions.class); dfOptions.setProject("project_name"); dfOptions.setStagingLocation("gs://bucket_name/staging"); // 修正路径 dfOptions.setTempLocation("gs://bucket_name/folder_name"); // 加载带权限范围的凭证 GoogleCredentials credentials = GoogleCredentials.fromStream(new FileInputStream(PATH_TO_JSON_KEY)) .createScoped(Lists.newArrayList("https://www.googleapis.com/auth/cloud-platform")); dfOptions.setGcpCredential(credentials); Pipeline myPipe= Pipeline.create(dfOptions);
方式二:使用服务账号邮箱(Dataflow Runner优先)
使用Dataflow Runner时,推荐通过serviceAccount指定服务账号,同时通过环境变量GOOGLE_APPLICATION_CREDENTIALS指定密钥文件,无需手动设置gcpCredential:
DataflowPipelineOptions dfOptions = PipelineOptionsFactory.as(DataflowPipelineOptions.class); dfOptions.setProject("project_name"); dfOptions.setStagingLocation("gs://bucket_name/staging"); dfOptions.setTempLocation("gs://bucket_name/folder_name"); dfOptions.setRunner(DataflowRunner.class); dfOptions.setServiceAccount("correct-service-account@project-id.iam.gserviceaccount.com"); // 修正拼写 Pipeline myPipe= Pipeline.create(dfOptions);
3. 解决依赖冲突
google-cloud-dataflow-java-sdk-all 2.5.0版本过旧,与Apache Beam 2.24存在严重依赖冲突(Dataflow SDK 2.5.0对应Beam 2.5,版本差异极大)。必须移除该依赖,因为从Beam 2.10开始,Dataflow Runner已集成在Apache Beam核心依赖中,无需单独引入旧版Dataflow SDK。
调整后保留核心依赖:
- Apache Beam 2.24相关依赖(如
beam-sdks-java-core、beam-runners-google-cloud-dataflow-java) - google-cloud-storage 2.11.3
- google-api-client 1.35.1
4. 本地运行的凭证优先级
本地运行时,Beam优先读取GOOGLE_APPLICATION_CREDENTIALS环境变量指定的密钥文件,可先设置该变量再运行管道:
export GOOGLE_APPLICATION_CREDENTIALS=/path/to/your/service-account-key.json
内容的提问来源于stack exchange,提问作者Gestalt
相关产品推荐
相关产品推荐

