如何在sparklyr中按时间窗口及不同行的值查找数据对?
在sparklyr中基于时间窗口和反向关系查找数据对
我理解你要找的是互为反向的记录对(比如一条记录是from=X, to=Y,另一条是from=Y, to=X),并且这两条记录的时间戳落在指定的时间窗口内。下面我会一步步带你实现这个需求,用你提供的测试数据来演示。
步骤1:准备数据并连接Spark
首先我们先把你的测试数据加载好,然后连接到Spark并转换为Spark DataFrame:
# 加载测试数据 elemuid <- c(1, 2, 3, 4, 5, 6, 7) timestamp <- c("2018-02-10 23:00:00", "2018-02-10 23:01:00", "2018-02-10 22:59:00", "2018-02-10 22:40:00", "2018-02-10 22:39:00", "2018-02-10 22:37:00", "2018-02-10 23:01:00") from <- c(10, 8, 2, 12, 7, 8, 9) to <- c(9, 10, 10, 3, 12, 7, 8) value <- c(56, 26, 60, 50, 90, 80, 50) df <- data.frame(elemuid, timestamp, from, to, value) # 加载sparklyr并连接Spark library(sparklyr) sc <- spark_connect(master = "local") # 本地模式,可根据你的运行环境调整 spark_df <- copy_to(sc, df, "test_data", overwrite = TRUE)
步骤2:转换时间戳类型
Spark需要时间字段是timestamp类型才能进行时间差计算,所以先做类型转换:
spark_df <- spark_df %>% mutate(timestamp = to_timestamp(timestamp, "yyyy-MM-dd HH:mm:ss"))
步骤3:自连接实现反向配对+时间窗口过滤
核心思路是自连接(将表和自身连接),通过连接条件和过滤规则筛选符合要求的配对:
- 两条记录是反向关系:
a.from = b.to且a.to = b.from - 排除同一条记录的自配对:
a.elemuid != b.elemuid - 时间戳落在指定窗口内(这里设置为5分钟,你可以按需调整)
- 避免重复配对(比如(1,2)和(2,1)只保留一个)
# 设置时间窗口为5分钟(对应300秒) time_window_secs <- 300 matched_pairs <- spark_df %>% alias("a") %>% inner_join(spark_df %>% alias("b"), join_condition = expr("a.from = b.to AND a.to = b.from AND a.elemuid != b.elemuid")) %>% # 过滤时间差在窗口内的记录 filter(abs(unix_timestamp(a.timestamp) - unix_timestamp(b.timestamp)) <= time_window_secs) %>% # 避免重复配对(仅保留elemuid更小的在前的组合) filter(a.elemuid < b.elemuid) %>% # 选择并命名需要展示的字段,区分两个配对记录 select( elemuid_a = a.elemuid, timestamp_a = a.timestamp, from_a = a.from, to_a = a.to, value_a = a.value, elemuid_b = b.elemuid, timestamp_b = b.timestamp, from_b = b.from, to_b = b.to, value_b = b.value, time_diff_secs = abs(unix_timestamp(a.timestamp) - unix_timestamp(b.timestamp)) ) # 查看最终结果 collect(matched_pairs)
关键调整说明
- 时间窗口:如果需要调整窗口大小,直接修改
time_window_secs的值即可(比如10分钟就设为600)。 - 配对完整性:如果需要保留所有记录(包括找不到反向配对的),可以把
inner_join换成left_join。 - 重复配对:如果需要保留双向的配对结果(比如同时显示(1,2)和(2,1)),去掉
filter(a.elemuid < b.elemuid)这一行即可。
内容的提问来源于stack exchange,提问作者user60856839
相关产品推荐
相关产品推荐

