【问题标题】:pandas: apply filters taking into account timestamppandas:应用考虑时间戳的过滤器
【发布时间】:2020-04-02 19:20:45
【问题描述】:

我有以下测试数据:

import pandas as pd
import datetime

data = {'date': ['2014-01-01', '2014-01-02', '2014-01-03', '2014-01-04', '2014-01-05', '2014-01-06', '2014-01-07'],
     'id': [1, 2, 2, 3, 4, 4, 5], 'name': ['Darren', 'Sabrina', 'Steve', 'Sean', 'Ray', 'Stef', 'Dany']}
data = pd.DataFrame(data)
data['date'] = pd.to_datetime(data['date'])

问题是:回溯 x 天(从每个条目查看),是否有超过 y 个不同的名称共享相同的 id?

这是我编写的代码。在我的示例中,我返回 x=2 天并检查至少两个共享相同 ID 的不同名称 (y=1)。如果至少存在两个不同的名称,我在列表“result_store”中保存 1,否则为 0。当然,在这个例子中,如果 i 小于 x,则返回 x 天是不可能的,但是这个小不准确对于我。

def rule(data, x=2, y=1):
    result_store = []
    for i in range(data.shape[0]):
        id = data['id'][i]
        end_time = data['date'][i]
        start_time = end_time-datetime.timedelta(days=x)
        time_frame = data[(data['date'] >= start_time) & (data['date'] <= end_time)]
        time_frame = time_frame.loc[time_frame['id'] == id]
        distinct_names = time_frame['name'].nunique()
        if distinct_names > y:
            result_store.append(1)
        else:
            result_store.append(0)

    return result_store

结果是

[0, 0, 1, 0, 0, 1, 0]

实际上,我有数千行,我的解决方案非常慢。我也尝试过使用 parmap 对索引进行并行化,但速度提升也不令人满意。有没有更有效的方法来做到这一点?也许通过使用 pyspark?

谢谢!

【问题讨论】:

    标签: python pandas datetime pyspark timestamp


    【解决方案1】:

    这适用于 spark2.4array_distinct 仅适用于 2.4)。我使用了您提供的 DataFrame,并且 spark 推断列日期为 TimestampType 类型。为了使我的 spark 代码正常工作,日期列 必须是 TimestampType 类型。 window 函数返回 2 天,基于 same id,并收集名称列表。如果不同名称的个数>1,则输入1,否则输入0。

    下面的代码使用rangeBetween(-(86400*2),Window.currentRow)这基本上意味着包含currentRow然后返回2天,所以如果当前行日期是3,它将包括 [3,2,1]。如果您只想要当前行日期和前 1 天,您可以将 86400*2 替换为 86400*1

    #If you can't use spark2.4 or get stuck, please leave a comment. 
    
    
    
    from pyspark.sql import functions as F
    from pyspark.sql.window import Window
    
    df=spark.createDataFrame(data)
    
    w=Window().partitionBy("id").orderBy((F.col("date")).cast("long")).rangeBetween(-(86400*2),Window.currentRow)
    df.withColumn("no_distinct", F.size(F.array_distinct(F.collect_list("name").over(w))))\
      .withColumn("no_distinct", F.when(F.col("no_distinct")>1, F.lit(1)).otherwise(F.lit(0)))\
      .orderBy(F.col("date")).show()
    
    +-------------------+---+-------+-----------+
    |               date| id|   name|no_distinct|
    +-------------------+---+-------+-----------+
    |2014-01-01 00:00:00|  1| Darren|          0|
    |2014-01-02 00:00:00|  2|Sabrina|          0|
    |2014-01-03 00:00:00|  2|  Steve|          1|
    |2014-01-04 00:00:00|  3|   Sean|          0|
    |2014-01-05 00:00:00|  4|    Ray|          0|
    |2014-01-06 00:00:00|  4|   Stef|          1|
    |2014-01-07 00:00:00|  5|   Dany|          0|
    +-------------------+---+-------+-----------+
    

    【讨论】:

    • 非常感谢您的快速响应!目前我正在尝试启动并运行 pyspark,但遇到了一些问题:“异常:Java 网关进程在发送其端口号之前已退出”明天我将与一位同事会面以解决问题。
    猜你喜欢
    • 1970-01-01
    • 2021-02-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-05-18
    • 1970-01-01
    相关资源
    最近更新 更多