Fra*_*ank 4 python apache-spark pyspark
我pyspark使用 AWS EMR 和 ~15 m4.large 核心来处理 50Gb 数据。
每行数据都包含一天中特定时间的一些信息。我使用以下for循环来提取和聚合每小时的信息。最后我是union数据,因为我希望将结果保存在一个csv 文件中。
# daily_df is a empty pyspark DataFrame
for hour in range(24):
hourly_df = df.filter(hourFilter("Time")).groupby("Animal").agg(mean("weights"), sum("is_male"))
daily_df = daily_df.union(hourly_df)
Run Code Online (Sandbox Code Playgroud)
据我所知,我必须执行以下操作才能强制pyspark.sql.Dataframe对象保存到 1 个 csv 文件(约 1Mb)而不是 100 多个文件:
daily_df.coalesce(1).write.csv("some_local.csv")
Run Code Online (Sandbox Code Playgroud)
似乎花了大约70分钟才能完成这个进度,我想知道我是否可以通过使用collect()类似的方法使其更快?
daily_df_pandas = daily_df.collect()
daily_df_pandas.to_csv("some_local.csv")
Run Code Online (Sandbox Code Playgroud)
coalesce(1)一般来说, 和都很collect糟糕,但预期输出大小约为 1MB,这并不重要。它根本不应该成为这里的瓶颈。
一个简单的改进是删除loop-> filter->union并执行单个聚合:
df.groupby(hour("Time"), col("Animal")).agg(mean("weights"), sum("is_male"))
Run Code Online (Sandbox Code Playgroud)
spark.sql.shuffle.partitions如果这还不够,那么这里的问题很可能是配置(如果您还没有这样做,最好的起点可能是调整)。
| 归档时间: |
|
| 查看次数: |
14865 次 |
| 最近记录: |