【问题标题】:Import pyspark dataframe from multiple S3 buckets, with a column denoting which bucket the entry came from从多个 S3 存储桶导入 pyspark 数据帧,其中有一列表示条目来自哪个存储桶
【发布时间】:2019-12-16 00:06:20
【问题描述】:

我有一个按日期分区的 S3 存储桶列表。第一个桶名为 2019-12-1,第二个桶名为 2019-12-2,依此类推。

这些存储桶中的每一个都存储我正在读入 pyspark 数据帧的镶木地板文件。从这些存储桶中的每一个生成的 pyspark 数据帧具有完全相同的架构。我想做的是遍历这些桶,并将所有这些 parquet 文件存储到一个 pyspark 数据框中,该数据框有一个日期列,表示数据框中每个条目实际来自哪个桶。

因为单独导入每个存储桶时生成的数据帧的架构有很多层深(即每一行包含结构数组的结构等),我想将所有存储桶组合成一个数据帧的唯一方法是拥有一个具有单个“日期”列的数据框。 “日期”列的每一行都将保存该日期对应的 S3 存储桶的内容。

我可以用这一行读取所有日期:

df = spark.read.parquet("s3://my_bucket/*")

我看到有人通过在这一行附加一个“withColumn”调用来创建一个“日期”列,但我不记得是如何实现的。

【问题讨论】:

    标签: amazon-s3 pyspark pyspark-dataframes


    【解决方案1】:

    使用input_file_name(),您可以从文件路径中提取S3存储桶名称:

    df.withColumn("dates", split(regexp_replace(input_file_name(), "s3://", ""), "/").getItem(0))\
      .show()
    

    我们拆分文件名并获取与存储桶名称对应的第一部分。

    这也可以使用正则表达式s3:\/\/(.+?)\/(.+) 来完成,第一组是存储桶名称:

    df.withColumn("dates", regexp_extract(input_file_name(), "s3:\/\/(.+?)\/(.+)", 1)).show()
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-07-23
      • 1970-01-01
      • 1970-01-01
      • 2019-11-07
      • 2020-05-29
      • 1970-01-01
      • 2021-11-01
      • 2017-11-23
      相关资源
      最近更新 更多