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,示例数据如下:
| timestamp | node_name | sensor_val |
|---|---|---|
| 2021-02-01 00:01:00 | x24n15 | 23.5 |
| 2021-02-01 00:01:00 | x24n15 | 23.5 |
| 2021-02-01 00:01:00 | x24n16 | 23.5 |
| 2021-02-01 00:01:00 | x24n17 | 23.5 |
| 2021-02-01 00:01:00 | x24n18 | 23.5 |
| 2021-02-01 00:01:00 | x24n19 | 23.5 |
| 2021-02-01 00:01:00 | x24n20 | 23.5 |
| 2021-02-01 00:01:01 | x24n15 | 23.5 |
| 2021-02-01 00:01:01 | x24n15 | 23.5 |
| 2021-02-01 00:01:01 | x24n16 | 23.5 |
| 2021-02-01 00:01:01 | x24n17 | 23.5 |
| 2021-02-01 00:01:01 | x24n18 | 23.5 |
| 2021-02-01 00:01:01 | x24n19 | 23.5 |
| 2021-02-01 00:01:01 | x24n20 | 23.5 |
业务需求
逐行匹配df1的规则,筛选df2中满足以下两个条件的记录:
timestamp落在对应行begin_time~end_time区间内node_name属于对应行node_name数组
关联上匹配到的df1的id字段,最终输出包含timestamp、node_name、sensor_val、id的结果集,预期结果示例如下:
| timestamp | node_name | sensor_val | id |
|---|---|---|---|
| 2021-02-01 00:01:00 | x24n15 | 23.5 | 1182 |
| 2021-02-01 00:01:00 | x24n15 | 23.5 | 1182 |
| 2021-02-01 00:01:00 | x24n16 | 23.5 | 1182 |
| 2021-02-01 00:01:00 | x24n17 | 23.5 | 1183 |
| 2021-02-01 00:01:00 | x24n18 | 23.5 | 1183 |
| 2021-02-01 00:01:00 | x24n19 | 23.5 | 1184 |
| 2021-02-01 00:01:00 | x24n20 | 23.5 | 1184 |
| 2021-02-01 00:01:01 | x24n15 | 23.5 | 1184 |
| 2021-02-01 00:01:01 | x24n15 | 23.5 | 1184 |
| 2021-02-01 00:01:01 | x24n16 | 23.5 | 1185 |
| 2021-02-01 00:01:01 | x24n17 | 23.5 | 1185 |
| 2021-02-01 00:01:01 | x24n18 | 23.5 | 1185 |
| 2021-02-01 00:01:01 | x24n19 | 23.5 | 1185 |
| 2021-02-01 00:01:01 | x24n20 | 23.5 | 1185 |
| 更多亿级规模数据行 |
现有实现问题
当前已尝试的实现方案执行速度极慢,无法在生产环境使用,代码如下:
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
相关产品推荐
相关产品推荐

