You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

不使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.28 14:54:03