【问题标题】:Within a loop logic to write CSV file - Spark DF, Equal Splits在循环逻辑内写入 CSV 文件 - Spark DF,Equal Splits
【发布时间】:2021-10-18 14:49:52
【问题描述】:

我使用本地 Hadoop Spark 实例。 我编写了下面的代码,将我的 spark df 分成 4 个部分。 我希望在此过程中将每个部分写入 CSV(即文件 1、文件 2 ..4) 我知道我可以使用df.toPandas().to_csv('file1.csv'),但我不确定如何获得 4 个单独的文件。

为什么?因为我的笔记本电脑内存问题。导出整个文件时出现 java 错误。如果我导出 10% 的数据就可以了。

# Define the number of splits you want
n_splits = 4
  
# Calculate count of each dataframe rows
each_len = df.count() // n_splits
  
# Create a copy of original dataframe
copy_df = df
  
# Iterate for each dataframe
i = 0
while i < n_splits:
  
    # Get the top `each_len` number of rows
    temp_df = copy_df.limit(each_len)
  
    # Truncate the `copy_df` to remove
    # the contents fetched for `temp_df`
    copy_df = copy_df.subtract(temp_df)
  
    # View the dataframe
    temp_df.show(truncate=False)
  
    # Increment the split number
    i += 1

【问题讨论】:

  • df.repartition(4).write.csv(...)
  • 这段代码给了我同样的“org.apache.spark.SparkException: Job aborted”。尝试将整个 DF 导出到 CSV 时出现错误。我需要一些解决方法来简化内存,即一步一步地,因为我一次只能导出 10% 的 DF。

标签: python apache-spark pyspark export-to-csv


【解决方案1】:

试试这个

var dataFile = “path/to/folder“ df.repartition(4).write.mode("overwrite").option("header","true").format("csv").save(dataFile)

我不确定本地实例,但 hadoop/mapreduce 旨在为每个 reduce 任务输出一个分区。如果您只在单个节点上运行,repart 可能会占用一些内存,这通常会给您 org.apache.spark.SparkException: Job aborted 错误。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-01-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-03-14
    相关资源
    最近更新 更多