【问题标题】:PySpark read csv from zip file in s3 with two different file typesPySpark 从 s3 中的 zip 文件中读取 csv,有两种不同的文件类型
【发布时间】:2021-08-12 22:51:28
【问题描述】:

我有一个包含 CSV 和 json 映射文件的 zip 文件。我想将 csv 读入 spark 数据框,并将 json 映射文件读入字典。我已经完成了后面的部分:

import boto3

obj = s3.get_object(Bucket='bucket', Key='key')

z = zipfile.ZipFile(io.BytesIO(obj["Body"].read()))

csvjson = json.loads(z.open(files[1]).read().decode('utf-8'))

一般来说,我想从 csv 文件中获取 df:

dfRaw = spark.read \
    .format("text") \
    .option("multiLine","true") \
    .option("inferSchema","false") \
    .option("header","true") \
    .option("ignoreLeadingWhiteSpace","true") \
    .option("ignoreTrailingWhiteSpace","true") \
    .load(z.open(files[0]).read().decode('utf-8'))

但这显然不起作用,因为load() 需要一个文件路径,而不是行本身。如何从 zip 文件中读取此文件到 spark 数据框中?

【问题讨论】:

  • 使用sc.parallelize(...) 加载它然后使用to_csv 怎么样?
  • @pltc 你能发布一个例子吗?我认为我在这里挂断的部分是从 zip 存档中访问它

标签: python apache-spark amazon-s3 pyspark


【解决方案1】:

由于您是手动“解压缩”CSV 文件并将输出作为字符串获取,因此您可以使用parallelize,如下所示

z = zipfile.ZipFile(io.BytesIO(obj["Body"].read()))
csv = [l.decode('utf-8').replace('\n', '') for l in z.open(files[0]).readlines()]

(spark
    .sparkContext
    .parallelize(csv)
    .toDF(T.StringType())
    .withColumn('value', F.from_csv('value', 'ID int, Trxn_Date string')) # your schema goes here
    .select('value.*')
    .show(10, False)
)

# Output
+----+----------+
|ID  |Trxn_Date |
+----+----------+
|null|Trxn_Date |
|100 |2021-03-24|
|133 |2021-01-22|
+----+----------+

【讨论】:

    猜你喜欢
    • 2012-03-09
    • 1970-01-01
    • 1970-01-01
    • 2018-07-17
    • 1970-01-01
    • 2021-07-05
    • 1970-01-01
    • 1970-01-01
    • 2021-10-25
    相关资源
    最近更新 更多