GCP Dataproc无法访问Cloud Storage的问题排查求助
问题描述
运行基于GCP Dataproc的Spark Structured Streaming作业,从Kafka拉取数据并处理,检查点存储在GCS存储桶gs://ss-checkpoint-10m-test-v1中,作业启动时报错无法访问存储桶内的metadata对象,返回403 Forbidden。
报错信息
Traceback (most recent call last): File "/tmp/3c930be06f4242378251aa6760390f8d/main-v1-test.py", line 352, in <module> sys.exit(main()) File "/tmp/3c930be06f4242378251aa6760390f8d/main-v1-test.py", line 348, in main ss.readData() File "/tmp/3c930be06f4242378251aa6760390f8d/main-v1-test.py", line 301, in readData query = df_stream.selectExpr("CAST(value AS STRING)", "timestamp", "topic").writeStream \ File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/streaming.py", line 1491, in start File "/usr/lib/spark/python/lib/py4j-0.10.9-src.zip/py4j/java_gateway.py", line 1304, in __call__ File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 111, in deco File "/usr/lib/spark/python/lib/py4j-0.10.9-src.zip/py4j/protocol.py", line 326, in get_return_value py4j.protocol.Py4JJavaError: An error occurred while calling o96.start. : java.io.IOException: Error accessing gs://ss-checkpoint-10m-test-v1/metadata at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageImpl.getObject(GoogleCloudStorageImpl.java:2221) at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageImpl.getItemInfo(GoogleCloudStorageImpl.java:2108) at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageFileSystem.getFileInfoInternal(GoogleCloudStorageFileSystem.java:1091) at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageFileSystem.getFileInfo(GoogleCloudStorageFileSystem.java:1065) at com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystemBase.getFileStatus(GoogleHadoopFileSystemBase.java:955) at org.apache.hadoop.fs.FileSystem.exists(FileSystem.java:1691) at org.apache.spark.sql.execution.streaming.StreamMetadata$.read(StreamMetadata.scala:53) at org.apache.spark.sql.execution.streaming.StreamExecution.<init>(StreamExecution.scala:174) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.<init>(MicroBatchExecution.scala:50) at org.apache.spark.sql.streaming.StreamingQueryManager.createQuery(StreamingQueryManager.scala:317) at org.apache.spark.sql.streaming.StreamingQueryManager.startQuery(StreamingQueryManager.scala:359) at org.apache.spark.sql.streaming.DataStreamWriter.startQuery(DataStreamWriter.scala:466) at org.apache.spark.sql.streaming.DataStreamWriter.startInternal(DataStreamWriter.scala:414) at org.apache.spark.sql.streaming.DataStreamWriter.start(DataStreamWriter.scala:301) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.GatewayConnection.run(GatewayConnection.java:238) at java.lang.Thread.run(Thread.java:750) Caused by: com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.googleapis.json.GoogleJsonResponseException: 403 Forbidden GET https://storage.googleapis.com/storage/v1/b/ss-checkpoint-10m-test-v1/o/metadata?fields=bucket,name,timeCreated,updated,generation,metageneration,size,contentType,contentEncoding,md5Hash,crc32c,metadata { "code" : 403, "errors" : [ { "domain" : "global", "message" : "939354532596-compute@developer.gserviceaccount.com does not have storage.objects.get access to the Google Cloud Storage object. Permission 'storage.objects.get' denied on resource (or it may not exist).", "reason" : "forbidden" } ], "message" : "939354532596-compute@developer.gserviceaccount.com does not have storage.objects.get access to the Google Cloud Storage object. Permission 'storage.objects.get' denied on resource (or it may not exist)." } at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.googleapis.json.GoogleJsonResponseException.from(GoogleJsonResponseException.java:146) at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:118) at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:37) at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.googleapis.services.AbstractGoogleClientRequest$1.interceptResponse(AbstractGoogleClientRequest.java:428) at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.http.HttpRequest.execute(HttpRequest.java:1111) at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:514) at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:455) at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.googleapis.services.AbstractGoogleClientRequest.execute(AbstractGoogleClientRequest.java:565) at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageImpl.getObject(GoogleCloudStorageImpl.java:2215) ... 24 more
附加信息
账号使用情况
- 本地启动作业时使用的认证账号:
kafka-admin@versa-sml-googl.iam.gserviceaccount.com - Dataproc虚拟机实际访问GCS的服务账号:
939354532596-compute@developer.gserviceaccount.com(项目默认Compute Engine服务账号)
本地认证账号列表
(base) Karans-MacBook-Pro:Versa-StructuredStreaming karanalang$ gcloud auth list Credentialed Accounts ACTIVE ACCOUNT dataproc-access@versa-kafka-poc.iam.gserviceaccount.com * kafka-admin@versa-sml-googl.iam.gserviceaccount.com karan.alang@gmail.com karan@versa-networks.com loki-748@versa-sml-googl.iam.gserviceaccount.com loki-storage@versa-sml-googl.iam.gserviceaccount.com
服务账号角色权限
kafka-admin@versa-sml-googl.iam.gserviceaccount.com
(base) Karans-MacBook-Pro:Versa-StructuredStreaming karanalang$ gcloud projects get-iam-policy versa-sml-googl \ > --flatten="bindings[].members" \ > --format='table(bindings.role)' \ > --filter="bindings.members:kafka-admin@versa-sml-googl.iam.gserviceaccount.com" ROLE roles/compute.admin roles/compute.storageAdmin roles/container.admin roles/container.clusterAdmin roles/dataproc.admin roles/dataproc.editor roles/dataproc.worker roles/iam.serviceAccountTokenCreator roles/iam.serviceAccountUser roles/metastore.editor roles/owner roles/storage.admin roles/storage.objectAdmin roles/storage.objectViewer
939354532596-compute@developer.gserviceaccount.com
(base) Karans-MacBook-Pro:Versa-StructuredStreaming karanalang$ gcloud projects get-iam-policy versa-sml-googl \ > --flatten="bindings[].members" \ > --format='table(bindings.role)' \ > --filter="939354532596-compute@developer.gserviceaccount.com" ROLE roles/compute.storageAdmin roles/editor roles/metastore.serviceAgent roles/storage.objectAdmin roles/storage.objectViewer
排查与修复建议
检查存储桶级IAM权限:项目级角色可能被存储桶的自定义IAM策略覆盖。运行以下命令查看存储桶的IAM配置:
gsutil iam get gs://ss-checkpoint-10m-test-v1确认
939354532596-compute@developer.gserviceaccount.com是否在桶的IAM列表中,且拥有storage.objects.get、storage.objects.list权限。若存储桶设置了拒绝规则,会优先于项目级权限生效。检查对象级ACL:如果
gs://ss-checkpoint-10m-test-v1/metadata对象已存在,可能单独配置了ACL。运行以下命令查看对象权限:gsutil acl get gs://ss-checkpoint-10m-test-v1/metadata确保服务账号对该对象有读取权限。
指定Dataproc集群服务账号:默认情况下,Dataproc集群使用项目的Compute Engine默认服务账号。重新创建集群时,可指定
kafka-admin@versa-sml-googl.iam.gserviceaccount.com作为集群服务账号,命令示例:gcloud dataproc clusters create <集群名称> \ --region <区域> \ --service-account kafka-admin@versa-sml-googl.iam.gserviceaccount.com \ --scopes cloud-platform验证权限生效状态:GCP IAM权限变更可能有5-10分钟的延迟,若刚修改过权限,等待一段时间后重新运行作业。
在Dataproc节点上测试权限:登录Dataproc集群的主节点,用默认服务账号直接测试访问存储桶:
gsutil ls gs://ss-checkpoint-10m-test-v1 gsutil cat gs://ss-checkpoint-10m-test-v1/metadata若测试失败,确认是权限问题;若测试成功,检查Spark作业的检查点路径配置是否正确,或Hadoop配置是否存在异常。
确认存储桶及对象存在性:运行
gsutil ls gs://ss-checkpoint-10m-test-v1确认存储桶存在,若metadata对象不存在,检查作业是否首次运行时因权限问题未创建成功。
内容的提问来源于stack exchange,提问作者Karan Alang

