Apply function after groupby in Sparklyr

Viewed 149

I have a sparklyr dataframe roughly like this

full_df %>% head()
absolute_time file_number battery_soc ...
2020-09-01 04:57:45 7 99.0 ...
2020-09-01 04:57:47 7 98.0 ...
2020-09-01 04:58:31 7 95.0 ...
2020-09-01 04:33:31 21 32.0 ...
2020-09-24 04:35:42 21 31.5 ...
2020-09-24 04:52:35 21 30.0 ...

I need to run an interpolation of the battery_soc column to smoothen its values for each file_number (while making sure the whole thing is ordered by absolute_time (an as.Date column).

The expected result would be like so:

absolute_time file_number battery_soc battery_interpolated ...
2020-09-01 04:57:45 7 99.0 99.0 ...
2020-09-01 04:57:47 7 98.0 97.5 ...
2020-09-01 04:58:31 7 95.0 95.0 ...
2020-09-01 04:33:31 21 32.0 32.0 ...
2020-09-24 04:35:42 21 31.5 31.0 ...
2020-09-24 04:52:35 21 30.0 30.0 ...

I'm trying the following:

full_df %>%
  select(file_number, battery_soc) %>%
  spark_apply(
    function(e) stats::smooth(e$battery_soc),
    names = "battery_interpolated",
    group_by = "file_number")

But that gives me an error:

Error : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 4554.0 failed 4 times, most recent failure: Lost task 0.3 in stage 4554.0 (TID 12296) (10.221.193.135 executor 4): java.lang.Exception: sparklyr worker rscript failure with status 255, check worker logs for details.
    at sparklyr.Rscript.init(rscript.scala:83)
    at sparklyr.WorkerApply$$anon$2.run(workerapply.scala:138)

Driver stacktrace:
    at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2783)
    at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2(DAGScheduler.scala:2730)
    at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2$adapted(DAGScheduler.scala:2724)
    at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62)
    at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55)
    at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49)
    at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:2724)
    at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1(DAGScheduler.scala:1260)
    at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1$adapted(DAGScheduler.scala:1260)
    at scala.Option.foreach(Option.scala:407)
    at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:1260)
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2991)
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2932)
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2920)
    at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
    at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:1033)
    at org.apache.spark.SparkContext.runJobInternal(SparkContext.scala:2476)
    at org.apache.spark.sql.execution.collect.Collector.runSparkJobs(Collector.scala:264)
    at org.apache.spark.sql.execution.collect.Collector.collect(Collector.scala:299)
    at org.apache.spark.sql.execution.collect.Collector$.collect(Collector.scala:82)
    at org.apache.spark.sql.execution.collect.Collector$.collect(Collector.scala:88)
    at org.apache.spark.sql.execution.collect.InternalRowFormat$.collect(cachedSparkResults.scala:75)
    at org.apache.spark.sql.execution.collect.InternalRowFormat$.collect(cachedSparkResults.scala:62)
    at org.apache.spark.sql.execution.ResultCacheManager.$anonfun$getOrComputeResultInternal$1(ResultCacheManager.scala:512)
    at scala.Option.getOrElse(Option.scala:189)
    at org.apache.spark.sql.execution.ResultCacheManager.getOrComputeResultInternal(ResultCacheManager.scala:511)
    at org.apache.spark.sql.execution.ResultCacheManager.getOrComputeResult(ResultCacheManager.scala:399)
    at org.apache.spark.sql.execution.ResultCacheManager.getOrComputeResult(ResultCacheManager.scala:374)
    at org.apache.spark.sql.execution.SparkPlan.executeCollectResult(SparkPlan.scala:406)
    at org.apache.spark.sql.Dataset.collectResult(Dataset.scala:3041)
    at org.apache.spark.sql.Dataset.collectFromPlan(Dataset.scala:3833)
    at org.apache.spark.sql.Dataset.$anonfun$collect$1(Dataset.scala:3008)
    at org.apache.spark.sql.Dataset.$anonfun$withAction$1(Dataset.scala:3825)
    at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withCustomExecutionEnv$5(SQLExecution.scala:130)
    at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:273)
    at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withCustomExecutionEnv$1(SQLExecution.scala:104)
    at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:854)
    at org.apache.spark.sql.execution.SQLExecution$.withCustomExecutionEnv(SQLExecution.scala:77)
    at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:223)
    at org.apache.spark.sql.Dataset.withAction(Dataset.scala:3823)
    at org.apache.spark.sql.Dataset.collect(Dataset.scala:3008)
    at sparklyr.Utils$.collect(utils.scala:26)
    at sparklyr.Utils.collect(utils.scala)
    at sun.reflect.GeneratedMethodAccessor522.invoke(Unknown Source)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at sparklyr.Invoke.invoke(invoke.scala:161)
    at sparklyr.StreamHandler.handleMethodCall(stream.scala:138)
    at sparklyr.StreamHandler.read(stream.scala:62)
    at sparklyr.BackendHandler.$anonfun$channelRead0$1(handler.scala:59)
    at scala.util.control.Breaks.breakable(Breaks.scala:42)
    at sparklyr.BackendHandler.channelRead0(handler.scala:40)
    at sparklyr.BackendHandler.channelRead0(handler.scala:14)
    at io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:99)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
    at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
    at io.netty.handler.codec.MessageToMessageDecoder.channelRead(MessageToMessageDecoder.java:103)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
    at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
    at io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:324)
    at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:296)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
    at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:357)
    at io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:379)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:365)
    at io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919)
    at io.netty.channel.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:163)
    at io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:714)
    at io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:650)
    at io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:576)
    at io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:493)
    at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:989)
    at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)
    at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
    at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.Exception: sparklyr worker rscript failure with status 255, check worker logs for details.
    at sparklyr.Rscript.init(rscript.scala:83)
    at sparklyr.WorkerApply$$anon$2.run(workerapply.scala:138)

Some(<code style = 'font-size:10p'> Error: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 4554.0 failed 4 times, most recent failure: Lost task 0.3 in stage 4554.0 (TID 12296) (10.221.193.135 executor 4): java.lang.Exception: sparklyr worker rscript failure with status 255, check worker logs for details. </code>)

Moreover, that doesn't ensure that my rows will be ordered by absolute_time, which is vital for this matter.

0 Answers
Related