不使用hadoop-aws如何通过Spark查询S3上的JSON文件
不依赖hadoop-aws实现Spark查询S3 JSON文件的方案
完全可以通过「分布式拉取S3文件内容 + 直接读取内存中JSON数据」的方式实现,不需要对接Hadoop文件系统抽象层,也不会引入hadoop-aws的性能问题,实现逻辑非常轻量。
具体实现步骤
- 引入依赖:仅需引入轻量的AWS S3官方SDK即可,不用整套引入hadoop-aws相关依赖
<!-- 以Maven依赖为例 --> <dependency> <groupId>com.amazonaws</groupId> <artifactId>aws-java-sdk-s3</artifactId> <version>1.12.XXX</version> </dependency>
- Driver端拉取待处理的S3 JSON文件列表,筛选出符合要求的对象key
- 把文件key分布式下发到Executor,每个分区初始化一次S3客户端拉取文件内容,生成JSON内容RDD
- 直接用Spark SQL的
read.json接口读取RDD生成DataFrame,后续正常执行SQL查询即可
代码示例
import com.amazonaws.services.s3.AmazonS3ClientBuilder import org.apache.spark.sql.SparkSession object S3JsonQueryDemo { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("S3JsonQueryWithoutHadoopAws") .getOrCreate() import spark.implicits._ // 替换为实际的桶名和文件前缀 val s3Bucket = "your-business-bucket" val jsonPathPrefix = "data/json/20240501/" // Driver端拉取目标JSON对象列表 val driverS3Client = AmazonS3ClientBuilder.defaultClient() import scala.collection.JavaConverters._ val targetJsonKeys = driverS3Client.listObjects(s3Bucket, jsonPathPrefix) .getObjectSummaries.asScala .map(_.getKey) .filter(_.endsWith(".json")) .toSeq // 分布式拉取JSON内容 val jsonRdd = spark.sparkContext.parallelize(targetJsonKeys, numSlices = Math.min(targetJsonKeys.size, 500)) .mapPartitions { keyIter => // 每个Executor分区初始化一次S3客户端,减少资源开销 val execS3Client = AmazonS3ClientBuilder.defaultClient() keyIter.map { key => val s3Obj = execS3Client.getObject(s3Bucket, key) // 大文件可以改为按行读取,避免全量加载内存溢出 scala.io.Source.fromInputStream(s3Obj.getObjectContent).mkString } } // 直接读取内存中的JSON内容生成DataFrame val jsonDf = spark.read.json(jsonRdd.toDS()) // 执行SQL查询逻辑 jsonDf.createOrReplaceTempView("tmp_s3_json") val resultDf = spark.sql("select count(*) as total, type from tmp_s3_json group by type") resultDf.show() } }
注意事项
- 内存适配:如果单JSON文件超过100M,建议改为按行读取S3对象流,不要直接全量读取成字符串,避免Executor OOM
- 权限配置:优先使用EC2实例角色、容器服务账户角色等无密钥方式配置S3访问权限,不要硬编码AK/SK
- 分区调整:根据总文件大小调整并行度,单个分区数据量控制在128M左右即可,拉取性能最优
这个方案完全不需要修改Spark底层的FileFormat、HadoopFSRelation等逻辑,开发量极小,也不会触发hadoop-aws处理parquet时的多余API调用问题,完全满足你在单个任务中完成S3文件下载+SQL查询的需求。
内容的提问来源于stack exchange,提问作者mridul_verma
相关产品推荐
相关产品推荐

