【问题标题】:Fill in null with previously known good value with pyspark使用 pyspark 用先前已知的良好值填充 null
【发布时间】:2016-07-20 12:02:46
【问题描述】:

有没有办法用最后一个有效值替换 pyspark 数据框中的 null 值?如果您认为 Windows 分区和排序需要它们,还有额外的 timestampsession 列。更具体地说,我想实现以下转换:

+---------+-----------+-----------+      +---------+-----------+-----------+
| session | timestamp |         id|      | session | timestamp |         id|
+---------+-----------+-----------+      +---------+-----------+-----------+
|        1|          1|       null|      |        1|          1|       null|
|        1|          2|        109|      |        1|          2|        109|
|        1|          3|       null|      |        1|          3|        109|
|        1|          4|       null|      |        1|          4|        109|
|        1|          5|        109| =>   |        1|          5|        109|
|        1|          6|       null|      |        1|          6|        109|
|        1|          7|        110|      |        1|          7|        110|
|        1|          8|       null|      |        1|          8|        110|
|        1|          9|       null|      |        1|          9|        110|
|        1|         10|       null|      |        1|         10|        110|
+---------+-----------+-----------+      +---------+-----------+-----------+

【问题讨论】:

  • 你不能。 DataFrames 行之间没有顺序。
  • 如果我有timestamp的订单怎么办?
  • 你不能按某种寡妇划分吗?这种情况下怎么办,手动一一处理条目并保持状态?
  • @eliasah “不可能”是一个强有力的断言,我会谨慎使用。正如下面的几个答案所证明的那样,它可能的。 (尽管这些解决方案可能并不适用于所有情况。)
  • @lostsoul29 我的评论是针对当时问题的状态给出的,现在已经过时了。我会删除它。谢谢!

标签: apache-spark pyspark apache-spark-sql


【解决方案1】:

我相信我有一个比公认的更简单的解决方案。它也使用函数,但使用名为“LAST”的函数并忽略空值。

让我们重新创建类似于原始数据的东西:

import sys
from pyspark.sql.window import Window
import pyspark.sql.functions as func

d = [{'session': 1, 'ts': 1}, {'session': 1, 'ts': 2, 'id': 109}, {'session': 1, 'ts': 3}, {'session': 1, 'ts': 4, 'id': 110}, {'session': 1, 'ts': 5},  {'session': 1, 'ts': 6}]
df = spark.createDataFrame(d)

打印出来:

+-------+---+----+
|session| ts|  id|
+-------+---+----+
|      1|  1|null|
|      1|  2| 109|
|      1|  3|null|
|      1|  4| 110|
|      1|  5|null|
|      1|  6|null|
+-------+---+----+

现在,如果我们使用窗口函数 LAST:

df.withColumn("id", func.last('id', True).over(Window.partitionBy('session').orderBy('ts').rowsBetween(-sys.maxsize, 0))).show()

我们刚刚得到:

+-------+---+----+
|session| ts|  id|
+-------+---+----+
|      1|  1|null|
|      1|  2| 109|
|      1|  3| 109|
|      1|  4| 110|
|      1|  5| 110|
|      1|  6| 110|
+-------+---+----+

希望对你有帮助!

【讨论】:

  • 请注意:此答案会将每个会话的所有行收集到某个执行程序节点。如果某些会话中的行数大于执行程序节点的内存,这将导致作业失败。
【解决方案2】:

@Oleksiy 的回答很好,但不能完全满足我的要求。在一个会话中,如果观察到多个nulls,则都用会话的第一个非null 填充。我需要 lastnull 值向前传播。

以下调整适用于我的用例:

def fill_forward(df, id_column, key_column, fill_column):

    # Fill null's with last *non null* value in the window
    ff = df.withColumn(
        'fill_fwd',
        func.last(fill_column, True) # True: fill with last non-null
        .over(
            Window.partitionBy(id_column)
            .orderBy(key_column)
            .rowsBetween(-sys.maxsize, 0))
        )

    # Drop the old column and rename the new column
    ff_out = ff.drop(fill_column).withColumnRenamed('fill_fwd', fill_column)

    return ff_out

【讨论】:

    【解决方案3】:

    这似乎是在使用Window functions

    import sys
    from pyspark.sql.window import Window
    import pyspark.sql.functions as func
    
    def fill_nulls(df):
        df_na = df.na.fill(-1)
        lag = df_na.withColumn('id_lag', func.lag('id', default=-1)\
                               .over(Window.partitionBy('session')\
                                     .orderBy('timestamp')))
    
        switch = lag.withColumn('id_change',
                                ((lag['id'] != lag['id_lag']) &
                                 (lag['id'] != -1)).cast('integer'))
    
    
        switch_sess = switch.withColumn(
            'sub_session',
            func.sum("id_change")
            .over(
                Window.partitionBy("session")
                .orderBy("timestamp")
                .rowsBetween(-sys.maxsize, 0))
        )
    
        fid = switch_sess.withColumn('nn_id',
                               func.first('id')\
                               .over(Window.partitionBy('session', 'sub_session')\
                                     .orderBy('timestamp')))
    
        fid_na = fid.replace(-1, 'null')
    
        ff = fid_na.drop('id').drop('id_lag')\
                              .drop('id_change')\
                              .drop('sub_session').\
                              withColumnRenamed('nn_id', 'id')
    
        return ff
    

    这是完整的null_test.py

    【讨论】:

    • @eliasah:你能复习一下答案吗?
    • 我现在真的在看。
    • 添加了测试,如果有帮助,可以让生活更轻松
    • 我正在写我的测试!谢谢。答案对我来说似乎很干净。拥有会话是必不可少的,这使得分区成为可能,从而使用窗口功能!
    • 不错的解决方案!我很惊讶 Spark 中还没有这个功能。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-03-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-11-23
    • 2010-11-08
    • 2020-03-25
    相关资源
    最近更新 更多