【问题标题】:How to write / writeStream each row of a dataframe into a different delta table如何将数据帧的每一行写入/写入流到不同的增量表中
【发布时间】:2019-06-28 17:52:26
【问题描述】:

我的数据框的每一行都有一个 CSV 内容。

我正在努力将每一行保存在不同的特定表中。

我认为我需要使用 foreach 或 UDF 来完成此操作,但这根本行不通。

我设法找到的所有内容就像 foreachs 中的简单打印或使用 .collect() 的代码(我真的不想使用)。

我也找到了重新分区的方式,但这不允许我选择每一行的去向。

rows = df.count()
df.repartition(rows).write.csv('save-dir')

你能给我一个简单而有效的例子吗?

【问题讨论】:

    标签: pyspark azure-databricks delta-lake


    【解决方案1】:

    将每一行保存为表格是一项昂贵的操作,不建议这样做。但是您正在尝试的可以像这样实现-

    df.write.format("delta").partitionBy("<primary-key-column>").save("/delta/save-dir")
    

    现在每一行都将保存为.parquet 格式,您可以从每个分区创建外部表。仅当您对每一行(即主键)都具有唯一值时,这才有效。

    【讨论】:

    • 我没有唯一键,实际上很多行都在同一个表中
    • 数据框有列 CSV | ID。我将使用 ID 来获取保存 CSV 的位置的信息。包括表名、数据库名、schema 和 sparkSchema。这就是为什么我需要一个 foreach 或 UDF。那是一切都失败的时候
    【解决方案2】:

    好吧,总而言之,它一如既往地非常简单,但我没有看到这一点。

    基本上,当您执行 foreach 并且您要保存的数据框构建在循环内时。 worker和driver不同,保存时不会自动设置“/dbfs/”路径,所以如果不手动添加“/dbfs/”,它会在worker本地保存数据。

    这就是我的循环不起作用的原因。

    【讨论】:

      【解决方案3】:

      你有没有试过.mode("append").repartionBy("ID"),它会为每个ID创建一个目录,然后别忘了放模式

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2023-03-14
        • 2020-11-03
        • 1970-01-01
        • 1970-01-01
        • 2021-07-13
        • 1970-01-01
        • 2022-07-05
        • 1970-01-01
        相关资源
        最近更新 更多