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

Spark Scala:基于近似位置与时间范围关联DataFrame并统计符合条件行数

我来帮你搞定这个需求!基于你给出的两个DataFrame结构和isNearUDF函数,这里有个清晰的Spark实现方案,不管你用Scala还是Python都能适配:

实现方案

1. 关联并过滤时间范围

首先我们把DF1和DF2做关联,同时先过滤掉DF2.DateTime不在DF1.StartDate与DF1.EndDate之间的无效数据,这样能减少后续计算的数据集大小:

Scala 代码

val joinedDF = DF1.join(DF2, DF2("DateTime").between(DF1("StartDate"), DF1("EndDate")))

Python 代码

joined_df = DF1.join(DF2, DF2["DateTime"].between(DF1["StartDate"], DF1["EndDate"]))

2. 筛选位置近似的记录

接着调用你已有的isNearUDF函数,过滤出两个DataFrame中位置满足近似条件的记录:

Scala 代码

val filteredDF = joinedDF.filter(isNearUDF(DF1("Position"), DF2("Position")))

Python 代码

filtered_df = joined_df.filter(isNearUDF(DF1["Position"], DF2["Position"]))

3. 按ID聚合统计行数

最后对过滤后的结果按DF1的ID分组,统计每个ID对应的符合条件的DF2行数,并重命名统计列让结果更清晰:

Scala 代码

val resultDF = filteredDF.groupBy(DF1("ID"))
  .count()
  .withColumnRenamed("count", "MatchedDF2Rows")

Python 代码

result_df = filtered_df.groupBy(DF1["ID"])
  .count()
  .withColumnRenamed("count", "MatchedDF2Rows")

优化小技巧

如果你的数据集规模较大,交叉关联可能会有性能瓶颈,你可以试试下面的优化方式:

  • 把时间过滤和位置过滤合并到join的条件中,一步完成关联+过滤,减少中间数据量:
// Scala 示例
val resultDF = DF1.join(DF2, 
  DF2("DateTime").between(DF1("StartDate"), DF1("EndDate")) && isNearUDF(DF1("Position"), DF2("Position"))
)
.groupBy(DF1("ID"))
.count()
.withColumnRenamed("count", "MatchedDF2Rows")
  • 给DF2的DateTime和Position字段添加分区或索引,提升关联查询的效率。

内容的提问来源于stack exchange,提问作者Nakeuh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 07:39:03