|
val _df_out = transformer.transform(_df) |
Сейчас ДАГ выполняется два раза - первый раз для получения левого дф для джойна, второй раз - для получения правого (который с примененными stats-функциями). Это должно решиться кэшированием на уровне спарка примерно так:
override def transform(_df: DataFrame): DataFrame = {
val cachedDf = _df.cache()
val _df_out = transformer.transform(cachedDf)
val res = positionalsMap.get("by") match {
case Some(Positional("by", List())) => cachedDf.drop(fieldsGenerated:_*).crossJoin(_df_out)
case Some(Positional("by", byList)) => cachedDf.drop(fieldsGenerated:_*).join(_df_out, byList.map(_.stripBackticks()))
case _ => _df_out
}
cachedDf.unpersist()
res
}
dispatcher/src/main/scala/ot/scalaotl/commands/OTLEventstats.scala
Line 17 in 3f3c877
Сейчас ДАГ выполняется два раза - первый раз для получения левого дф для джойна, второй раз - для получения правого (который с примененными stats-функциями). Это должно решиться кэшированием на уровне спарка примерно так: