【问题标题】:pyspark data pipeline use intermediary resultspyspark 数据管道使用中间结果
【发布时间】:2019-08-12 05:21:53
【问题描述】:

在 pyspark 中,我会对数据帧进行连续操作,并希望从中间结果中获取输出。不过,它总是需要同样的时间,我想知道它是否缓存过任何东西?换一种问法,使用中间结果的最佳做法是什么?在dask you can dodd.compute(df.amount.max(), df.amount.min()) 中将找出需要缓存和计算的内容。 pyspark 中是否有等价物?

在下面的例子中,当它到达print() 时会执行 3x 吗?

df_purchase = spark.read.parquet("s3a:/example/location")[['col1','col2']]
df_orders = df_purchase.groupby(['col1']).agg(pyspark.sql.functions.first("col2")).withColumnRenamed("first(col2, false)", "col2")
df_orders_clean = df_orders.dropna(subset=['col2'])

print(df_purchase.count(), df_orders.count(), df_orders_clean.count())

【问题讨论】:

    标签: pyspark


    【解决方案1】:

    是的,每次您对 dag 执行操作时。它执行并优化完整的查询。

    默认情况下,Spark 不缓存任何内容。

    缓存时要小心,缓存可能会以负面方式干扰:Spark: Explicit caching can interfere with Catalyst optimizer's ability to optimize some queries?

    【讨论】:

    • 罗杰谢谢!似乎如果我添加 df_orders.perist(); df_orders_clean.persist() 它在第一次运行时运行得更快,在后续运行中甚至更快。
    • @citynorman,我的荣幸!缓存时要小心,我更新了我的答案:)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-03-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-05
    • 1970-01-01
    相关资源
    最近更新 更多