使用Spark从S3复制到HDFS时遭遇ConnectionTimeOutException求助
我来帮你梳理这个连接超时的问题,结合你的配置和报错信息,整理几个关键的排查点和修复方案:
1. 替换旧的S3文件系统实现
你当前配置的org.apache.hadoop.fs.s3native.NativeS3FileSystem是Hadoop早期的S3实现,已经不再推荐使用,而且你设置的超时、代理等新参数对它完全不生效。这很可能是超时问题的核心原因,建议直接替换为官方推荐的S3A文件系统:
sparkSession.sparkContext.hadoopConfiguration.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")
2. 修正代理与SSL的配置冲突
你开启了SSL连接(fs.s3a.connection.ssl.enabled=true),但代理端口设置的是80——这通常是HTTP代理的端口,HTTPS代理一般用443或者其他指定端口。需要根据你的代理实际情况调整:
- 如果代理支持HTTPS转发,把
fs.s3a.proxy.port改为代理的HTTPS端口(比如443) - 如果代理仅支持HTTP,可临时关闭SSL(不推荐生产环境):
fs.s3a.connection.ssl.enabled=false
如果你的代理需要身份认证,补充以下配置:
sparkSession.sparkContext.hadoopConfiguration.set("fs.s3a.proxy.username", "your-proxy-username") sparkSession.sparkContext.hadoopConfiguration.set("fs.s3a.proxy.password", "your-proxy-password")
3. 开启路径风格的S3访问
报错里的目标地址是acron-avro-bucket.s3.amazonaws.com,这是虚拟主机风格的访问方式。部分代理服务器对这种带bucket名称的HTTPS主机名处理存在问题,建议开启路径风格访问:
sparkSession.sparkContext.hadoopConfiguration.set("fs.s3a.path.style.access", "true")
开启后,S3访问路径会变为s3.us-east-1.amazonaws.com/acron-avro-bucket/...,避免虚拟主机解析带来的代理兼容性问题。
4. 优化S3A的超时与重试参数
换用S3A后,补充更全面的超时和重试配置,提升连接稳定性:
// 延长socket超时时间 sparkSession.sparkContext.hadoopConfiguration.set("fs.s3a.socket.timeout", "60000") // 配置重试间隔和最大重试次数 sparkSession.sparkContext.hadoopConfiguration.set("fs.s3a.retry.interval", "2000") sparkSession.sparkContext.hadoopConfiguration.set("fs.s3a.retry.max.tries", "10")
5. 修复AVRO文件读取方式
注意到你用textFile读取AVRO文件,这是错误的——textFile仅适用于纯文本文件,AVRO是二进制格式,应该使用Spark的AVRO数据源:
// 确保已添加spark-avro依赖 val df = sparkSession.read.avro("s3a://acrXXXXXXXXXXXXXXXXX5.avro") // 若需保存为ObjectFile,转换为RDD后操作 df.rdd.saveAsObjectFile("hdfs://c411apy.int.westgroup.com:8020/project/ecpdevingest/avro/100")
虽然这不是超时的直接原因,但错误的读取方式可能引发额外异常,干扰问题排查。
额外优化:Spark参数配置方式
你的Spark资源参数(比如spark.executor.instances)应该通过SparkConf设置,而非hadoopConfiguration,确保参数生效:
val sparkConf = new SparkConf() .setMaster("yarn") .setAppName("To hdfs") .set("spark.executor.instances", "8") .set("spark.executor.cores", "4") .set("spark.executor.memory", "32g") .set("spark.driver.memory", "4g") .set("spark.driver.cores", "2") .set("spark.yarn.queue","root.ecpdevingest") .set("spark.speculation", "false") val sparkSession = SparkSession.builder().config(sparkConf).getOrCreate()
按上述步骤调整后,重新提交作业应该能解决连接超时问题。
内容的提问来源于stack exchange,提问作者Swathi

