如何为Spark数据集补全日期缺失的accountid对应0值记录?
当然可以实现!这个需求本质上是要生成日期维度和账号维度的笛卡尔积,然后和原数据做左关联,把缺失的val填充为0。下面我给你一步步拆解实现过程,用你提供的示例数据来演示:
实现步骤
1. 提取所有唯一的日期和账号
首先我们需要拿到数据里所有不重复的日期和账号,这样才能生成完整的维度组合:
// 提取所有唯一日期 val dates = df.select("date").distinct() // 提取所有唯一账号 val accountIds = df.select("accountid").distinct()
2. 生成日期和账号的笛卡尔积
用crossJoin方法得到所有可能的日期-账号组合,这一步就是补全缺失的组合:
val fullCombination = dates.crossJoin(accountIds)
3. 和原数据左关联并填充缺失值
把完整组合和原数据集左关联,然后用coalesce函数把null的val替换成0.0:
import org.apache.spark.sql.functions._ val resultDF = fullCombination.join(df, Seq("date", "accountid"), "left") .select( col("date"), coalesce(col("val"), lit(0.0)).alias("val"), col("accountid") )
4. 查看结果
执行resultDF.show()就能得到你想要的输出了:
+----------+-----+---------+ | date| val|accountid| +----------+-----+---------+ |2018-01-01|100.5| id1| |2018-01-01| 0.0| id2| |2018-01-02|120.6| id1| |2018-01-02| 0.0| id2| |2018-01-03| 0.0| id1| |2018-01-03|450.2| id2| +----------+-----+---------+
完整代码整合
把上面的步骤整合起来,完整可运行的代码如下:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ object FillMissingRecords { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("FillMissingRecords") .master("local[*]") .getOrCreate() import spark.implicits._ // 原数据集 val df = spark.sparkContext.parallelize(Seq( ("2018-01-01", 100.5,"id1"), ("2018-01-02", 120.6,"id1"), ("2018-01-03", 450.2,"id2") )).toDF("date", "val","accountid") // 提取唯一日期和账号 val dates = df.select("date").distinct() val accountIds = df.select("accountid").distinct() // 生成完整组合并左关联填充0 val resultDF = dates.crossJoin(accountIds) .join(df, Seq("date", "accountid"), "left") .select( col("date"), coalesce(col("val"), lit(0.0)).alias("val"), col("accountid") ) // 展示结果 resultDF.show() } }
补充说明
- 如果你的日期维度不是来自原数据(比如需要覆盖一个指定的时间范围),可以用
sequence函数生成连续日期序列,再和账号做笛卡尔积,示例代码如下:
// 生成2018-01-01到2018-01-03的日期序列 val dateRange = spark.sql("select sequence(to_date('2018-01-01'), to_date('2018-01-03'), interval 1 day) as date") .select(explode(col("date")).alias("date"))
- 注意数据类型匹配:原数据的
val是Double类型,所以用lit(0.0)而不是lit(0),避免出现类型不兼容的问题。
内容的提问来源于stack exchange,提问作者Masterbuilder
相关产品推荐
相关产品推荐

