You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 06:56:42