【问题标题】:Azure Databricks writing a file into Azure Data Lake Gen 2Azure Databricks 将文件写入 Azure Data Lake Gen 2
【发布时间】:2019-09-23 10:55:33
【问题描述】:

我有一个 Azure Data Lake gen1 和一个 Azure Data Lake gen2(带分层结构的 Blob 存储),我正在尝试创建一个 Databricks 笔记本 (Scala),它读取 2 个文件并将一个新文件写回 Data Lake。在 Gen1 和 Gen2 中,我都遇到了同样的问题,我指定的输出 csv 的文件名被保存为一个目录,并且在该目录中它正在写入 4 个文件“committed, started 、_SUCCESS 和 part-00000-tid-

对于我的生活,我无法弄清楚为什么它会这样做并且实际上没有将 csv 保存到该位置。 这是我编写的代码示例。如果我在 df_join 数据帧上执行 .show() ,那么它会输出正确的结果。但是 .write 不能正常工作。

val df_names = spark.read.option("header", "true").csv("/mnt/datalake/raw/names.csv")
val df_addresses = spark.read.option("header", "true").csv("/mnt/datalake/raw/addresses.csv")

val df_join = df_names.join(df_addresses, df_names.col("pk") === df_addresses.col("namepk"))


df_join.write
.format("com.databricks.spark.csv")
.option("header", "true")
.mode("overwrite")
.save("/mnt/datalake/reports/testoutput.csv")

【问题讨论】:

    标签: scala azure azure-data-lake databricks azure-databricks


    【解决方案1】:

    如果我正确理解您的需求,您只想将 Spark DataFrame 数据写入 Azure Data Lake 中名为 testoutput.csv 的单个 csv 文件,而不是带有一些分区文件的名为 testoutput.csv 的目录。

    所以你不能直接通过使用DataFrameWriter.save这样的Spark函数来实现它,因为实际上dataframe writer是基于Azure Data Lake将数据写入HDFS。 HDFS 将数据保存为一个名为 yours 的目录和一些分区文件。请参阅The Hadoop FileSystem API Definition等有关HDFS的一些文档了解它。

    然后,根据我的经验,您可以尝试在 Scala 程序中使用 Azure Data Lake SDK for Jave,将数据作为单个文件直接从 DataFrame 写入 Azure Data Lake。并且你可以参考一些示例https://github.com/Azure-Samples?utf8=%E2%9C%93&q=data-lake&type=&language=java。

    【讨论】:

      【解决方案2】:

      它创建包含多个文件的目录的原因是因为每个分区都单独保存并写入数据湖。要保存单个输出文件,您需要重新分区数据框

      让我们使用数据框 API

      confKey = "fs.azure.account.key.srcAcctName.blob.core.windows.net"
      secretKey = "==" #your secret key
      spark.conf.set(confKey,secretKey)
      blobUrl = 'wasbs://MyContainerName@srcAcctName.blob.core.windows.net'
      

      合并您的数据框

      df_join.coalesce(1)
      .write
      .format("com.databricks.spark.csv")
      .option("header", "true")
      .mode("overwrite")
      .save("blobUrl" + "/reports/")
      

      更改文件名

      files = dbutils.fs.ls(blobUrl + '/reports/')
      output_file = [x for x in files if x.name.startswith("part-")]
      dbutils.fs.mv(output_file[0].path, "%s/reports/testoutput.csv" % (blobUrl))
      

      【讨论】:

      • 感谢您的评论。这会将一个 csv 文件保存到数据湖中,但不会以指定的名称保存。它仍然忽略文件名,只写入文件为“part-00000-tid ....”的目录名。在做了一些进一步的研究之后,我认为不可能使用特定的文件名保存到数据湖中。这对我来说仍然很奇怪。
      • 如果你打算走这条路,你应该使用.coalesce(1)而不是.repartition(1)来减少节点之间的数据移动。 Repartition 将对所有节点的数据进行完全洗牌,然后再将其减少为单个分区。
      • @DavidP 你是 100% 正确的。我会相应地改变答案。
      【解决方案3】:

      试试这个:

      df_join.to_csv('/dbfs/mnt/....../df.csv', sep=',', header=True, index=False)
      

      【讨论】:

        猜你喜欢
        • 2020-11-22
        • 2022-11-10
        • 2023-03-04
        • 1970-01-01
        • 2020-01-24
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多