【问题标题】:Spark Persist Method is not getting invokedSpark Persist 方法没有被调用
【发布时间】:2020-06-05 18:12:45
【问题描述】:

我无法理解persist(StorageLevel.Memory_and_disk) 的工作原理。 我有下面的代码 sn-p 工作正常。

val df1= Spark.read.from.hive.table()
            // Perform a high complex calculation & Aggregation.
            .toDF()
df1.write.toAnotherHiveTable.

140GB 的数据大约需要 1 小时。 阶段/任务有

1.CalcStage = 750tasks(耗时 50-55 分钟) 2.插入 Hive 阶段 = 100 个任务(3-4 分钟)

我有新的要求修改如下。

val df1= Spark.read.from.hive.table()
            // Perform a high complex calculation & Aggregation.
            .toDF()


val df2=df1.filter($"exchange" === "commodities")

val finalDF = df1.join(df2)
finalDF .write.toAnotherHiveTable.

对于相同数量的数据,这开始需要大约 1 小时 40 分钟。 并且 Stage/Task 有

1.CalcStage = 750tasks(耗时 1 小时 30 分钟) 2.CalcStage = 750tasks(taking 1 hr 30 mins) // 前 2 个阶段开始并行运行。两个阶段日志 已从 Hive 表条目中读取 3.插入 Hive 阶段 = 100 个任务(8-10 分钟)

我假设由于 df1 和 df2 依赖于 df1 计算逻辑,它会进行计算和聚合。我添加了 df1 的持久化如下。

val df1= Spark.read.from.hive.table()
            // Perform a high complex calculation & Aggregation.
            .toDF().persist(StorageLevel.Memory_and_disk)


val df2=df1.filter($"exchange" === "commodities")

val finalDF = df1.join(df2)
finalDF .write.toAnotherHiveTable

我认为添加持久性将帮助我减少运行完整计算/聚合的第二个 df。但是我错了。 DAG 计划和阶段日志与之前的运行相同。没有观察到变化。

我在这里遗漏了什么吗?请帮助我理解为什么坚持方法没有改变任何东西。

【问题讨论】:

  • 你能澄清你的“//执行一个高度复杂的计算和聚合。”?

标签: scala apache-spark apache-spark-sql persistence cloudera


【解决方案1】:

你需要调用一个动作 n 然后持久化你的第一个数据帧,然后计算才会保存在你的内存中

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-10-20
    • 2012-03-13
    • 2021-07-07
    • 2014-06-02
    • 2011-08-28
    • 2017-02-24
    • 2015-02-23
    相关资源
    最近更新 更多