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