【发布时间】: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