【发布时间】: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