【问题标题】:Optimal way to save spark sql dataframe to S3 using information stored in them使用存储在其中的信息将 spark sql 数据帧保存到 S3 的最佳方法
【发布时间】:2019-05-28 21:21:54
【问题描述】:

我的数据框包含如下数据:

        channel  eventId1               eventId2               eventTs  eventTs2  serialNumber  someCode
        Web-DTB akefTEdZhXt8EqzLKXNt1Wjg    akTEdZhXt8EqzLKXNt1Wjg  1545502751154   1545502766731   4   rfs
        Web-DTB 3ycLHHrbEkBJ.piYNyI7u55w    3ycLHHEkBJ.piYNyI7u55w  1545502766247   1545502767800   4   njs
        Web-DTB 3ycL4rHHEkBJ.piYNyI7u55w    3ycLHHEkBJ.piYNyI7u55w  1545502766247   1545502767800   4   null

我需要将此数据保存到 S3 路径,如下所示:

  s3://test/data/ABC/hb/eventTs/[eventTs]/uploadTime_[eventTs2]/*.json.gz

我需要如何从分区中提取数据以写入 S3 路径:(s3 路径是数据帧中存在的 eventTs 和 eventTs2 的函数)

df.write.partitionBy("eventTs","eventTs2").format("json").save("s3://test/data/ABC/hb????")

我想我可以遍历数据框中的每一行,提取路径并保存到 S3,但不想这样做。

有没有办法按 eventTs 和 eventTs2 上的数据帧分组,然后将数据帧保存到完整的 S3 路径?有什么更优化的吗?

【问题讨论】:

    标签: scala apache-spark amazon-s3 apache-spark-sql


    【解决方案1】:

    Spark 支持类似于 Hive 中的分区。如果 eventTs, eventTs2 的不同元素数量较少,分区将是解决此问题的好方法。

    查看scala doc 了解有关 partitionBy 的更多信息。

    示例用法:

    val someDF = Seq((1, "bat", "marvel"), (2, "mouse", "disney"), (3, "horse", "animal"), (1, "batman", "marvel"), (2, "tom", "disney") ).toDF("id", "name", "place")
    someDF.write.partitionBy("id", "name").orc("/tmp/somedf")
    

    如果您在“id”和“name”上写入带有 paritionBy 的数据框,则会创建以下目录结构。

    /tmp/somedf/id=1/name=bat
    /tmp/somedf/id=1/name=batman
    
    /tmp/somedf/id=2/name=mouse
    /tmp/somedf/id=2/name=tom
    
    /tmp/somedf/id=3/name=horse
    

    第一个和第二个分区变成目录,所有id等于1,name为bat的行都会保存在目录结构/tmp/somedf/id=1/name=bat下,partitionBy中定义的分区顺序决定了目录的顺序。

    在您的情况下,分区将位于 eventTs 和 eventTS2 上。

    val someDF = Seq(
            ("Web-DTB","akefTEdZhXt8EqzLKXNt1Wjg","akTEdZhXt8EqzLKXNt1Wjg","1545502751154","1545502766731",4,"rfs"),
            ("Web-DTB","3ycLHHrbEkBJ.piYNyI7u55w","3ycLHHEkBJ.piYNyI7u55w","1545502766247","1545502767800",4,"njs"),
            ("Web-DTB","3ycL4rHHEkBJ.piYNyI7u55w","3ycLHHEkBJ.piYNyI7u55w","1545502766247","1545502767800",4,"null"))
        .toDF("channel" , "eventId1", "eventId2", "eventTs",  "eventTs2",  "serialNumber",  "someCode")
    someDF.write("eventTs", "eventTs2").orc("/tmp/someDF")
    

    如下创建目录结构。

    /tmp/someDF/eventTs=1545502766247/eventTs2=1545502767800
    /tmp/someDF/eventTs=1545502751154/eventTs2=1545502766731
    

    【讨论】:

    • 除了我正在考虑存储在 S3 中。我知道这种分区逻辑。在 S3 中寻找一种简单而干净的方法。
    • 分区将是最简单的恕我直言,除非 eventTS/eventTS2 中不同元素的数量不是数千个,并且数据框中没有数千个分区。在这些情况下,您最终会为每个分区创建数千个非常小的文件。
    • 编辑了问题,以便更清楚地了解我卡在哪里。 S3 路径是 eventTs 和 eventTs2 的函数,我怀疑我们是否可以在 S3 中保存而不指定您在 HDFS 中存储的完整路径。
    • 你可以在不指定完整路径的情况下保存在 S3 中,试试看。
    • 有效!不知道为什么我认为它不会。无论如何在基于此进行分区后删除一列?还是在 partitionBy 中使用 UDF?因为我不想格式化日期并从纪元以 yyyyMMdd 格式存储新列并增加数据大小。而且我不能删除原始列的粒度。因为我有 3 个与原始数据列不同的分区列。
    猜你喜欢
    • 2018-05-20
    • 2015-10-07
    • 2012-06-19
    • 2019-01-30
    • 2021-10-07
    • 1970-01-01
    • 2015-07-21
    • 1970-01-01
    • 2021-07-23
    相关资源
    最近更新 更多