使用Spark 1.4.0,Scala 2.10
我一直试图找出一种方法来使用最后一次已知的观察来转发填充空值,但我没有看到一种简单的方法.我认为这是一件非常常见的事情,但找不到显示如何执行此操作的示例.
我看到函数向前转移填充NaN的值,或滞后/超前函数来填充或移位数据偏移量,但没有任何东西可以获取最后的已知值.
在线查看,我在R中看到很多关于同一件事的Q/A,但在Spark/Scala中没有.
我正在考虑在日期范围内进行映射,从结果中过滤出NaN并选择最后一个元素,但我想我对语法感到困惑.
使用DataFrames我尝试类似的东西
import org.apache.spark.sql.expressions.Window
val sqlContext = new HiveContext(sc)
var spec = Window.orderBy("Date")
val df = sqlContext.read.format("com.databricks.spark.csv").option("header", "true").option("inferSchema", "true").load("test.csv")
val df2 = df.withColumn("testForwardFill", (90 to 0).map(i=>lag(df.col("myValue"),i,0).over(spec)).filter(p=>p.getItem.isNotNull).last)
Run Code Online (Sandbox Code Playgroud)
但这并没有让我任何地方.
过滤器部分不起作用; map函数返回一个spark.sql.Columns序列,但是filter函数需要返回一个Boolean,所以我需要从Column中获取一个值来测试,但似乎只有Column方法返回一个Column.
有没有办法在Spark上更"简单"地做到这一点?
感谢您的输入
编辑:
简单示例示例输入:
2015-06-01,33
2015-06-02,
2015-06-03,
2015-06-04,
2015-06-05,22
2015-06-06,
2015-06-07,
...
Run Code Online (Sandbox Code Playgroud)
预期产量:
2015-06-01,33
2015-06-02,33
2015-06-03,33
2015-06-04,33
2015-06-05,22
2015-06-06,22
2015-06-07,22
Run Code Online (Sandbox Code Playgroud)
注意:
编辑:
按照@ zero323的回答我试过这样:
import org.apache.spark.sql.Row
import org.apache.spark.rdd.RDD
val rows: RDD[Row] = df.orderBy($"Date").rdd
def notMissing(row: Row): Boolean = { !row.isNullAt(1) } …Run Code Online (Sandbox Code Playgroud) 使用Spark 1.5.1,
我一直在尝试用我的DataFrame的一列的最后一个已知观测值来填充空值。
可以从空值开始,在这种情况下,我将使用第一个已知的观察向后填充该空值。但是,如果这也使代码复杂化,则可以跳过这一点。
在这篇文章中,zero323提供了一个针对Scala的解决方案,用于解决非常相似的问题。
但是,我不了解Scala,也无法在Pyspark API代码中“翻译”它。可以用Pyspark做到吗?
谢谢你的帮助。
下面是一个简单的示例输入示例:
| cookie_ID | Time | User_ID
| ------------- | -------- |-------------
| 1 | 2015-12-01 | null
| 1 | 2015-12-02 | U1
| 1 | 2015-12-03 | U1
| 1 | 2015-12-04 | null
| 1 | 2015-12-05 | null
| 1 | 2015-12-06 | U2
| 1 | 2015-12-07 | null
| 1 | 2015-12-08 | U1
| 1 …Run Code Online (Sandbox Code Playgroud) 遇到错误,我认为是由窗口函数引起的。
当我应用这个脚本并只保留几个示例行时,它工作正常但是当我将它应用到我的整个数据集(只有几 GB)时,它在最后一步尝试坚持到 hdfs 时失败,出现这个奇怪的错误......当我坚持不使用窗口函数时脚本工作,所以问题一定来自那个(我有大约 325 个特征列通过 for 循环运行)。
知道什么可能导致问题吗?我的目标是通过正向填充方法对数据框中的每个变量进行时间序列数据的估算。
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql import types as T
from pyspark.sql import Window
import sys
print(spark.version)
'2.3.0'
# sample data
df = spark.createDataFrame([('2019-05-10 7:30:05', '1', '10', '0.5', 'FALSE'),\
('2019-05-10 7:30:10', '2', 'UNKNOWN', '0.24', 'FALSE'),\
('2019-05-10 7:30:15', '3', '6', 'UNKNOWN', 'TRUE'),\
('2019-05-10 7:30:20', '4', '7', 'UNKNOWN', 'UNKNOWN'),\
('2019-05-10 7:30:25', '5', '10', '1.1', 'UNKNOWN'),\
('2019-05-10 7:30:30', '6', 'UNKNOWN', '1.1', 'NULL'),\
('2019-05-10 7:30:35', '7', …Run Code Online (Sandbox Code Playgroud)