【问题标题】:PySpark: how to groupby, resample and forward-fill null values?PySpark:如何分组、重新采样和前向填充空值?
【发布时间】:2019-08-17 04:46:44
【问题描述】:

考虑到以下数据集在 Spark 中,我想以特定频率(例如 5 分钟)重新采样日期。

START_DATE = dt.datetime(2019,8,15,20,33,0)
test_df = pd.DataFrame({
    'school_id': ['remote','remote','remote','remote','onsite','onsite','onsite','onsite','remote','remote'],
    'class_id': ['green', 'green', 'red', 'red', 'green', 'green', 'green', 'green', 'red', 'green'],
    'user_id': [15,15,16,16,15,17,17,17,16,17],
    'status': [0,1,1,1,0,1,0,1,1,0],
    'start': pd.date_range(start=START_DATE, periods=10, freq='2min')
})

test_df.groupby(['school_id', 'class_id', 'user_id', 'start']).min()

但是,我还希望在两个特定日期范围之间进行重新采样:2019-08-15 20:30:002019-08-15 21:00:00。因此,school_idclass_iduser_id 的每组将有 6 个条目,两个日期范围之间每 5 分钟一个桶。 重采样生成的null 条目应由前向填充填充。

我已将 Pandas 用于示例数据集,但实际数据帧将在 Spark 中提取,因此我正在寻找的方法也应在 Spark 中完成。

我猜这种方法可能类似于PySpark: how to resample frequencies,但我无法让它在这种情况下工作。

感谢您的帮助

【问题讨论】:

    标签: python pyspark


    【解决方案1】:

    这可能不是获得最终结果的最佳方式,只是想在这里展示一下想法。

    1. 首先,创建 DataFrame 并将时间戳转换为整数
    from datetime import datetime
    import pytz
    from pytz import timezone
    
    # Create DataFrame
    START_DATE = datetime(2019,8,15,20,33,0)
    test_df = pd.DataFrame({
        'school_id': ['remote','remote','remote','remote','onsite','onsite','onsite','onsite','remote','remote'],
        'class_id': ['green', 'green', 'red', 'red', 'green', 'green', 'green', 'green', 'red', 'green'],
        'user_id': [15,15,16,16,15,17,17,17,16,17],
        'status': [0,1,1,1,0,1,0,1,1,0],
        'start': pd.date_range(start=START_DATE, periods=10, freq='2min')
    })
    
    # Convert TimeStamp to Integers
    df = spark.createDataFrame(test_df)
    print(df.dtypes)
    df = df.withColumn('start', F.col('start').cast("bigint"))
    df.show()
    

    这个输出:

    +---------+--------+-------+------+----------+
    |school_id|class_id|user_id|status|     start|
    +---------+--------+-------+------+----------+
    |   remote|   green|     15|     0|1565915580|
    |   remote|   green|     15|     1|1565915700|
    |   remote|     red|     16|     1|1565915820|
    |   remote|     red|     16|     1|1565915940|
    |   onsite|   green|     15|     0|1565916060|
    |   onsite|   green|     17|     1|1565916180|
    |   onsite|   green|     17|     0|1565916300|
    |   onsite|   green|     17|     1|1565916420|
    |   remote|     red|     16|     1|1565916540|
    |   remote|   green|     17|     0|1565916660|
    +---------+--------+-------+------+----------+
    
    1. 创建您想要的时间序列
    # Create time sequece needed
    start = datetime.strptime('2019-08-15 20:30:00', '%Y-%m-%d %H:%M:%S')
    eastern = timezone('US/Eastern')
    start = eastern.localize(start)
    times = pd.date_range(start = start, periods = 6, freq='5min')
    times = [s.timestamp() for s in times]
    print(times)
    
    [1565915400.0, 1565915700.0, 1565916000.0, 1565916300.0, 1565916600.0, 1565916900.0]
    
    1. 最后,为每个组创建数据框
    # Use pandas_udf to create final DataFrame
    schm = StructType(df.schema.fields + [StructField('epoch', IntegerType(), True)])
    @pandas_udf(schm, PandasUDFType.GROUPED_MAP)
    def resample(pdf):
        pddf = pd.DataFrame({'epoch':times})
        pddf['school_id'] = pdf['school_id'][0]
        pddf['class_id'] = pdf['class_id'][0]
        pddf['user_id'] = pdf['user_id'][0]
    
    
        res = np.searchsorted(times, pdf['start'])
        arr = np.zeros(len(times))
        arr[:] = np.nan
        arr[res] = pdf['start']
        pddf['status'] = arr
    
        arr[:] = np.nan
        arr[res] = pdf['status']
        pddf['start'] = arr
        return pddf
    
    df = df.groupBy('school_id', 'class_id', 'user_id').apply(resample)
    df = df.withColumn('timestamp', F.to_timestamp(df['epoch']))
    df.show(60)
    
    

    最终结果:

    +---------+--------+-------+----------+-----+----------+-------------------+
    |school_id|class_id|user_id|    status|start|     epoch|          timestamp|
    +---------+--------+-------+----------+-----+----------+-------------------+
    |   remote|     red|     16|      null| null|1565915400|2019-08-15 20:30:00|
    |   remote|     red|     16|      null| null|1565915700|2019-08-15 20:35:00|
    |   remote|     red|     16|1565915940|    1|1565916000|2019-08-15 20:40:00|
    |   remote|     red|     16|      null| null|1565916300|2019-08-15 20:45:00|
    |   remote|     red|     16|1565916540|    1|1565916600|2019-08-15 20:50:00|
    |   remote|     red|     16|      null| null|1565916900|2019-08-15 20:55:00|
    |   onsite|   green|     15|      null| null|1565915400|2019-08-15 20:30:00|
    |   onsite|   green|     15|      null| null|1565915700|2019-08-15 20:35:00|
    |   onsite|   green|     15|      null| null|1565916000|2019-08-15 20:40:00|
    |   onsite|   green|     15|1565916060|    0|1565916300|2019-08-15 20:45:00|
    |   onsite|   green|     15|      null| null|1565916600|2019-08-15 20:50:00|
    |   onsite|   green|     15|      null| null|1565916900|2019-08-15 20:55:00|
    |   remote|   green|     17|      null| null|1565915400|2019-08-15 20:30:00|
    |   remote|   green|     17|      null| null|1565915700|2019-08-15 20:35:00|
    |   remote|   green|     17|      null| null|1565916000|2019-08-15 20:40:00|
    |   remote|   green|     17|      null| null|1565916300|2019-08-15 20:45:00|
    |   remote|   green|     17|      null| null|1565916600|2019-08-15 20:50:00|
    |   remote|   green|     17|1565916660|    0|1565916900|2019-08-15 20:55:00|
    |   onsite|   green|     17|      null| null|1565915400|2019-08-15 20:30:00|
    |   onsite|   green|     17|      null| null|1565915700|2019-08-15 20:35:00|
    |   onsite|   green|     17|      null| null|1565916000|2019-08-15 20:40:00|
    |   onsite|   green|     17|1565916180|    1|1565916300|2019-08-15 20:45:00|
    |   onsite|   green|     17|1565916420|    1|1565916600|2019-08-15 20:50:00|
    |   onsite|   green|     17|      null| null|1565916900|2019-08-15 20:55:00|
    |   remote|   green|     15|      null| null|1565915400|2019-08-15 20:30:00|
    |   remote|   green|     15|1565915580|    0|1565915700|2019-08-15 20:35:00|
    |   remote|   green|     15|      null| null|1565916000|2019-08-15 20:40:00|
    |   remote|   green|     15|      null| null|1565916300|2019-08-15 20:45:00|
    |   remote|   green|     15|      null| null|1565916600|2019-08-15 20:50:00|
    |   remote|   green|     15|      null| null|1565916900|2019-08-15 20:55:00|
    +---------+--------+-------+----------+-----+----------+-------------------+
    
    

    现在每个组都有 6 个时间戳。 注意,并不是所有的原始 'status' 和 'start' 都映射到最终的 DataFrame,这是因为在 resample udf 中,它发生在 5minute 区间,两个 'start' 时间可以映射到同一个时间网格点,你在这里失去一个。这可以在udf 中根据您的频率以及您希望如何保留数据进行调整。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-09-11
      • 1970-01-01
      • 2017-01-09
      • 1970-01-01
      • 1970-01-01
      • 2017-01-19
      • 1970-01-01
      相关资源
      最近更新 更多