AWS Glue ETL任务以Data Catalog表为输入时出现连接拒绝错误
AWS Glue ETL调用Data Catalog表时Spark连接拒绝问题修复
问题现象
运行以Glue Data Catalog表(底层数据存储在S3中)为输入的Glue ETL任务时,抛出Spark实例连接拒绝错误;直接以S3桶为输入(不使用Data Catalog表)时,任务可以正常运行。
错误日志
2021-11-09 03:09:26,992 ERROR [main] glue.ProcessLauncher (Logging.scala:logError(91)): Exception in User Class java.lang.reflect.UndeclaredThrowableException at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1862) at org.apache.spark.deploy.SparkHadoopUtil.runAsSparkUser(SparkHadoopUtil.scala:64) at org.apache.spark.executor.CoarseGrainedExecutorBackend$.run(CoarseGrainedExecutorBackend.scala:188) at org.apache.spark.executor.CoarseGrainedExecutorBackend$.main(CoarseGrainedExecutorBackend.scala:281) at org.apache.spark.executor.CoarseGrainedExecutorBackendPlugin$class.launch(CoarseGrainedExecutorBackendWrapper.scala:10) at org.apache.spark.executor.CoarseGrainedExecutorBackendWrapper$$anon$1.launch(CoarseGrainedExecutorBackendWrapper.scala:15) at org.apache.spark.executor.CoarseGrainedExecutorBackendWrapper.launch(CoarseGrainedExecutorBackendWrapper.scala:19) at org.apache.spark.executor.CoarseGrainedExecutorBackendWrapper$.main(CoarseGrainedExecutorBackendWrapper.scala:5) at org.apache.spark.executor.CoarseGrainedExecutorBackendWrapper.main(CoarseGrainedExecutorBackendWrapper.scala) 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 com.amazonaws.services.glue.SparkProcessLauncherPlugin$class.invoke(ProcessLauncher.scala:48) at com.amazonaws.services.glue.ProcessLauncher$$anon$1.invoke(ProcessLauncher.scala:78) at com.amazonaws.services.glue.ProcessLauncher.launch(ProcessLauncher.scala:133) at com.amazonaws.services.glue.ProcessLauncher$.main(ProcessLauncher.scala:30) at com.amazonaws.services.glue.ProcessLauncher.main(ProcessLauncher.scala) Caused by: org.apache.spark.SparkException: Exception thrown in awaitResult: at org.apache.spark.util.ThreadUtils$.awaitResult(ThreadUtils.scala:226) at org.apache.spark.rpc.RpcTimeout.awaitResult(RpcTimeout.scala:75) at org.apache.spark.rpc.RpcEnv.setupEndpointRefByURI(RpcEnv.scala:101) at org.apache.spark.executor.CoarseGrainedExecutorBackend$$anonfun$run$1.apply$mcV$sp(CoarseGrainedExecutorBackend.scala:201) at org.apache.spark.deploy.SparkHadoopUtil$$anon$2.run(SparkHadoopUtil.scala:65) at org.apache.spark.deploy.SparkHadoopUtil$$anon$2.run(SparkHadoopUtil.scala:64) at java.security.AccessController.doPrivileged(Native Method) at javax.security.auth.Subject.doAs(Subject.java:422) at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1844) ... 17 more Caused by: java.io.IOException: Failed to connect to /172.35.144.51:41533 at org.apache.spark.network.client.TransportClientFactory.createClient(TransportClientFactory.java:245) at org.apache.spark.network.client.TransportClientFactory.createClient(TransportClientFactory.java:187) at org.apache.spark.rpc.netty.NettyRpcEnv.createClient(NettyRpcEnv.scala:198) at org.apache.spark.rpc.netty.Outbox$$anon$1.call(Outbox.scala:194) at org.apache.spark.rpc.netty.Outbox$$anon$1.call(Outbox.scala:190) at java.util.concurrent.FutureTask.run(FutureTask.java:266) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:748) Caused by: io.netty.channel.AbstractChannel$AnnotatedConnectException: Connection refused: /172.35.144.51:41533 at sun.nio.ch.SocketChannelImpl.checkConnect(Native Method) at sun.nio.ch.SocketChannelImpl.finishConnect(SocketChannelImpl.java:716) at io.netty.channel.socket.nio.NioSocketChannel.doFinishConnect(NioSocketChannel.java:323) at io.netty.channel.nio.AbstractNioChannel$AbstractNioUnsafe.finishConnect(AbstractNioChannel.java:340) at io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:633) at io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:580) at io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:497) at io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:459) at io.netty.util.concurrent.SingleThreadEventExecutor$5.run(SingleThreadEventExecutor.java:858) at io.netty.util.concurrent.DefaultThreadFactory$DefaultRunnableDecorator.run(DefaultThreadFactory.java:138) ... 1 more Caused by: java.net.ConnectException: Connection refused ... 11 more
关联ETL脚本
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job args = getResolvedOptions(sys.argv, ["JOB_NAME"]) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args["JOB_NAME"], args) # Script generated for node Data Catalog table DataCatalogtable_node1 = glueContext.create_dynamic_frame.from_catalog( database="<PLACEHOLDERFORCATALOGDB>", table_name="<PLACEHOLDERFORCATALOGTABLE>", transformation_ctx="DataCatalogtable_node1", ) # Script generated for node ApplyMapping ApplyMapping_node2 = ApplyMapping.apply( frame=DataCatalogtable_node1, mappings=[ ("id", "long", "id", "long"), ("country", "string", "country", "string"), ("state", "string", "state", "string"), ("city", "string", "city", "string"), ("amount", "double", "amount", "double"), ], transformation_ctx="ApplyMapping_node2", ) # Script generated for node S3 bucket S3bucket_node3 = glueContext.write_dynamic_frame.from_options( frame=ApplyMapping_node2, connection_type="s3", format="csv", connection_options={"path": "s3://<PLACEHOLDERFORS3BUCKET>/", "partitionKeys": []}, transformation_ctx="S3bucket_node3", ) job.commit()
修复方案
该报错的核心原因是Spark Executor和Driver的RPC通信被拦截,或Data Catalog访问的网络、权限配置异常,直接读S3不需要拉取Data Catalog元数据,因此可以正常运行。按以下顺序排查修复即可:
- 检查Glue任务关联的安全组出站规则:允许TCP协议下全部端口(0-65535)的VPC内部IP段通信,Spark RPC通信会随机使用高位端口,默认安全组未放开内部互访时会触发连接拒绝
- 检查Glue任务的IAM角色权限:需附加
AWSGlueServiceRole托管策略,或手动添加glue:GetDatabase、glue:GetTable、glue:GetPartitions的Data Catalog访问权限 - 若Glue任务部署在VPC内,需配置Glue Data Catalog的VPC端点:确保VPC端点的安全组允许来自Glue任务所在安全组的443端口入站访问,同时开启VPC的DNS解析功能
- 可选配置优化:在Glue任务的作业参数中添加
--conf spark.network.timeout=600s,延长Spark网络超时时间,避免元数据拉取量大时触发的连接超时
内容的提问来源于stack exchange,提问作者Van
相关产品推荐
相关产品推荐

