代码之家  ›  专栏  ›  技术社区  ›  Georg Heiler

熊猫用不同的频率排列不规则的时间序列

  •  2
  • Georg Heiler  · 技术社区  · 7 年前

    ID 一天又一天。

    但是我有两个问题:

    1. ID, time, value
    2. 构造连接条件,即假设 LEFT

    编辑

    一些虚拟数据

    import pandas as pd
    from datetime import datetime
    import numpy as np
    def make_df(frequency, valueName):
        date_rng = pd.date_range(start='2018-01-01', end='2018-01-02', freq=frequency)
        ts = pd.Series(np.random.randn(len(date_rng)), index=date_rng)
        groups = ['a', 'b', 'c', 'd', 'e']
        group_series = [groups[np.random.randint(len(groups))] for i in range(0, len(date_rng))]
        df = pd.DataFrame(ts, columns=[valueName])
        df['group'] = group_series
        return df
    df_1 = make_df('ms', 'value_A')
    display(df_1.head())
    df_2 = make_df('H', 'value_B')
    display(df_2.head())
    df_3 = make_df('S', 'value_C')
    display(df_3.head())
    

    代码(不是真正的pythonic): 我在尝试一些类似于 a JOIN b ON a.group = b.group AND time in window(some_seconds) 在SQL中,如果有多个匹配的记录,即不仅第一个而且所有记录都匹配/生成一行,那么这就有问题了。

    此外,我还将数据分组为(spark): df.groupBy($"KEY", window($"time", "5 minutes")).sum("metric") 但这可能是相当有损失的。

    然后我发现(熊猫) Pandas aligning multiple dataframes with TimeStamp index 它看起来已经很有趣了,但是只产生精确的匹配。但是,当尝试使用 df_2.join(df_3, how='outer', on=['group'], rsuffix='_1') 它不仅在时间上,而且 group 它失败了,错误是 pd.concat

    https://github.com/twosigma/flint 它在一个时间间隔内实现了一个时间序列联接-但是,我在使用它时遇到了问题。

    1 回复  |  直到 7 年前
        1
  •  1
  •   Georg Heiler    7 年前

    弗林特是我选择的工具。起初,flint不会在spark 2.2上工作,但我在这里修复了: https://github.com/geoHeil/flint/commit/a2827d38e155ec8ddd4252dc62d89181f14f0c47 以下方法效果很好:

    val left = Seq((1,1L, 0.1), (1, 2L,0.2), (3,1L,0.3), (3, 2L,0.4)).toDF("groupA", "time", "valueA")
      val right = Seq((1,1L, 11), (1, 2L,12), (3,1L,13), (3, 2L,14)).toDF("groupB", "time", "valueB")
      val leftTs = TimeSeriesRDD.fromDF(dataFrame = left)(isSorted = false, timeUnit = MILLISECONDS)
      val rightTS        = TimeSeriesRDD.fromDF(dataFrame = right)(isSorted = false, timeUnit = MILLISECONDS)
    
      val mergedPerGroup = leftTs.leftJoin(rightTS, tolerance = "1s")
    

    mergedPerGroup.toDF.filter(col("groupA") === col("groupB")).show
    +-------+------+------+------+------+
    |   time|groupA|valueA|groupB|valueB|
    +-------+------+------+------+------+
    |1000000|     3|   0.3|     3|    13|
    |2000000|     3|   0.4|     3|    14|
    

    要删除重复项,请使用distinct。