【问题标题】:How to read huge number of .gz S3 files into RDD?如何将大量 .gz S3 文件读入 RDD?
【发布时间】:2020-09-29 23:52:24
【问题描述】:
aws s3api list-objects-v2 --bucket cw-milenko-tests | grep 'tick_c'

输出显示

    "Key": "Json_gzips/tick_calculated_3_2020-05-27T11-50-22.json.gz",
    "Key": "Json_gzips/tick_calculated_3_2020-05-27T11-52-59.json.gz",
    "Key": "Json_gzips/tick_calculated_3_2020-05-27T11-55-08.json.gz",
    "Key": "Json_gzips/tick_calculated_3_2020-05-27T11-57-30.json.gz",
    "Key": "Json_gzips/tick_calculated_3_2020-05-27T11-59-59.json.gz",
    "Key": "Json_gzips/tick_calculated_4_2020-05-27T09-14-28.json.gz",
    "Key": "Json_gzips/tick_calculated_4_2020-05-27T11-35-38.json.gz",  

与wc -l

aws s3api list-objects-v2 --bucket cw-milenko-tests | grep 'tick_c' | wc -l
457

我可以将一个文件读入数据框。

val path ="tick_calculated_2_2020-05-27T00-01-21.json"
scala> val tick1DF = spark.read.json(path)
tick1DF: org.apache.spark.sql.DataFrame = [aml_barcode_canc: string, aml_barcode_payoff: string ... 70 more fields]

我很惊讶地看到反对票。 我想知道的是如何将457个文件加载到RDD中?我看到this SO 的问题。 有可能吗?有什么限制? 这是我到目前为止所尝试的。

val rdd1 = sc.textFile("s3://cw-milenko-tests/Json_gzips/tick_calculated*.gz")

如果我选择 s3a

val rdd1 = sc.textFile("s3a://cw-milenko-tests/Json_gzips/tick_calculated*.gz")
rdd1: org.apache.spark.rdd.RDD[String] = s3a://cw-milenko-tests/Json_gzips/tick_calculated*.gz MapPartitionsRDD[3] at textFile at <console>:27

也不行。

尝试检查我的 RDD。

scala> rdd1.take(1)
java.io.IOException: No FileSystem for scheme: s3
  at org.apache.hadoop.fs.FileSystem.getFileSystemClass(FileSystem.java:2660)
  at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:2667)
  at org.apache.hadoop.fs.FileSystem.access$200(FileSystem.java:94)

文件系统未被识别。

我的目标:

s3://json.gz -> rdd -> parquet

【问题讨论】:

  • 就像 Spark 中的任何其他 JSON 文件一样,我会说?你读过文档吗? spark.apache.org/docs/latest/sql-data-sources-json.html - 实现起来非常简单,文档也很好,应该可以帮助到你
  • @UninformedUser 看看我的新编辑,我解释了我真正想要实现的目标。

标签: apache-spark pyspark


【解决方案1】:

试试这个-

  /**
      * /Json_gzips
      * |-  spark-test-data1.json.gz
      * --------------------
      * {"id":1,"name":"abc1"}
      * {"id":2,"name":"abc2"}
      * {"id":3,"name":"abc3"}
      */

    /**/Json_gzips
      *|-   spark-test-data2.json.gz
      * --------------------
      * {"id":1,"name":"abc1"}
      * {"id":2,"name":"abc2"}
      * {"id":3,"name":"abc3"}
      */
    val path = getClass.getResource("/Json_gzips").getPath
    // path till the root directory which contains the all .gz files
    spark.read.json(path).show(false)

    /**
      * +---+----+
      * |id |name|
      * +---+----+
      * |1  |abc1|
      * |2  |abc2|
      * |3  |abc3|
      * |1  |abc1|
      * |2  |abc2|
      * |3  |abc3|
      * +---+----+
      */

如果需要,您可以将此 df 转换为 rdd

【讨论】:

    【解决方案2】:
    from pyspark.sql import SparkSession
    
    //Create Spark Session
    spark = SparkSession 
        .builder 
        .appName("Python Spark SQL basic example") 
        .getOrCreate()
    
    //To read all files inside from S3 in under Json_gzips key
    df = spark.read.json("s3a://cw-milenko-tests/Json_gzips/tick_calculated*.gz")
    df.show()
    rdd = df.rdd // to convert it to rdd
    

    使用s3a 代替s3

    why s3a over s3?

    还为 hadoop-aws 2.7.3 和 AWS SDK 添加依赖项 Add AWS S3 supporting JARs

    【讨论】:

    • 请使用 s3a,因为这是您可以在 Spark 中访问 s3 文件的方式。 @miki_cloud
    • 您可能需要将 hadoop-aws 2.7.3 和 AWS sdk jar 添加到类路径
    猜你喜欢
    • 1970-01-01
    • 2015-02-13
    • 2017-02-16
    • 2019-09-05
    • 1970-01-01
    • 2020-02-12
    • 2014-07-24
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多