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

PySpark基于另一DataFrame的时间范围与节点条件过滤大表

PySpark 亿级数据区间+数组匹配关联性能问题

数据集说明

  • df1:共10000行,包含id、begin_time、end_time、node_name(数组类型)4个字段,示例数据如下:
+-------------+-------------------+-------------------+----------------------------------------------------------------------------
|      id     |begin_time         |end_time           |node_name                                                                  |
+-------------+-------------------+-------------------+---------------------------------------------------------------------------+
|   1182      |2021-02-01 00:01:00|2021-02-01 00:02:20|[x24n15, x1117, x1108, x1113, x11n02, x11n03, x11n11, x11n32]              |
|   1183      |2021-02-01 00:02:40|2021-02-01 00:03:50|[x28n02, x1112, x1109, x1110]                                              |
|   1184      |2021-02-01 00:04:10|2021-02-01 00:07:10|[x32n10, x34n13, x13n16, x32n09, x28n01]                                   |
|   1185      |2021-02-01 00:05:00|2021-02-01 00:06:30|[x50n09, x50n08]                                                           |
|   1186      |2021-02-01 00:07:00|2021-02-01 00:08:20|[x50n08]                                                                   |
  • df2:共5亿行,包含约100个字段,核心字段为timestamp、node_name、sensor_val,示例数据如下:
timestampnode_namesensor_val
2021-02-01 00:01:00x24n1523.5
2021-02-01 00:01:00x24n1523.5
2021-02-01 00:01:00x24n1623.5
2021-02-01 00:01:00x24n1723.5
2021-02-01 00:01:00x24n1823.5
2021-02-01 00:01:00x24n1923.5
2021-02-01 00:01:00x24n2023.5
2021-02-01 00:01:01x24n1523.5
2021-02-01 00:01:01x24n1523.5
2021-02-01 00:01:01x24n1623.5
2021-02-01 00:01:01x24n1723.5
2021-02-01 00:01:01x24n1823.5
2021-02-01 00:01:01x24n1923.5
2021-02-01 00:01:01x24n2023.5

业务需求

逐行匹配df1的规则,筛选df2中满足以下两个条件的记录:

  1. timestamp落在对应行begin_time~end_time区间内
  2. node_name属于对应行node_name数组

关联上匹配到的df1的id字段,最终输出包含timestamp、node_name、sensor_val、id的结果集,预期结果示例如下:

timestampnode_namesensor_valid
2021-02-01 00:01:00x24n1523.51182
2021-02-01 00:01:00x24n1523.51182
2021-02-01 00:01:00x24n1623.51182
2021-02-01 00:01:00x24n1723.51183
2021-02-01 00:01:00x24n1823.51183
2021-02-01 00:01:00x24n1923.51184
2021-02-01 00:01:00x24n2023.51184
2021-02-01 00:01:01x24n1523.51184
2021-02-01 00:01:01x24n1523.51184
2021-02-01 00:01:01x24n1623.51185
2021-02-01 00:01:01x24n1723.51185
2021-02-01 00:01:01x24n1823.51185
2021-02-01 00:01:01x24n1923.51185
2021-02-01 00:01:01x24n2023.51185
更多亿级规模数据行

现有实现问题

当前已尝试的实现方案执行速度极慢,无法在生产环境使用,代码如下:

SeriesAppend = []
data_collect= df1.rdd.toLocalIterator()
for row in data_collect:
    bt = row.begin_time
    et = row.end_time 
    temp_row_df = spark.sql("SELECT timestamp_t, node_name, sensor_val FROM df2_table WHERE timestamp_t >= '2021-02-01 00:00:00' AND timestamp_t < '2021-02-02 00:00:00' AND node_name IN row.node_name ") 
    temp_row_df = temp_row_df.withColumn("node_name", F.lit(row.allocation_id))
    SeriesAppend.append(temp_row_df)
df_series = reduce(DataFrame.unionAll, SeriesAppend)

寻求可落地的高性能实现方案。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.02 07:21:38