PySpark内连接性能优化:比原生SQL慢数个数量级的问题
PySpark内连接查询性能优化紧急求助
我在PySpark中实现简单内连接,相同逻辑的数据库原生SQL查询仅需约45秒,但PySpark版本已运行14小时仍未结束,预估至少需要5天才能完成,急需优化该查询。
详细场景
我有两张表PARENT和CHILD:
PARENT包含parent_id、timestamp及元数据列CHILD包含child_id、parent_id、value及元数据列
需求是对最近3小时的数据,基于parent_id执行内连接。数据规模:
PARENT总计约3亿行,过滤最近3小时后仅约10万行CHILD近10亿行
耗时45秒的原生SQL如下:
SELECT COUNT(*) FROM PARENT INNER JOIN CHILD ON PARENT.parent_id = CHILD.parent_id WHERE PARENT.timestamp >= SYSDATE - 3/24;
对应的初始PySpark代码(原本期望和原生SQL耗时相近,结果完全不符合预期):
from datetime import datetime, timedelta from pyspark.sql import functions as F # 加载PARENT到parent,CHILD到child # ... parent = parent.where(F.col("timestamp") >= datetime.now() - timedelta(hours=3)) joined = parent.join(child, "parent_id")
尝试过的优化方案
通过Spark UI查看DAG后发现,慢查询的核心原因是Spark会扫描CHILD的全部10亿行。于是我尝试先提取过滤后PARENT的parent_id集合,再用这个集合过滤CHILD,修改后的代码如下:
from datetime import datetime, timedelta from pyspark.sql import functions as F # 加载PARENT到parent,CHILD到child # ... parent = parent.where(F.col("timestamp") >= datetime.now() - timedelta(hours=3)) parent_id_rows = parent.select("parent_id").collect() parent_ids = [row.parent_id for row in parent_id_rows] child = child.where(F.col("parent_id").isin(parent_ids)) joined = parent.join(child, "parent_id")
当前困境
修改后的脚本仍存在相同问题——Spark依然会扫描CHILD的全部10亿行。修改后脚本运行14小时后的执行计划显示,全表扫描CHILD的操作仍在执行。
求助问题
如何加速该查询?如何避免CHILD全表扫描?
运行时信息
- Java版本:21.0.2 (Private Build)
- Scala版本:2.12.18
Spark属性
- spark.default.parallelism: 256
- spark.driver.extraJavaOptions: -Djava.net.preferIPv6Addresses=false -XX:+IgnoreUnrecognizedVMOptions --add-opens=java.base/java.lang=ALL-UNNAMED --add-opens=java.base/java.lang.invoke=ALL-UNNAMED --add-opens=java.base/java.lang.reflect=ALL-UNNAMED --add-opens=java.base/java.io=ALL-UNNAMED --add-opens=java.base/java.net=ALL-UNNAMED --add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/java.util=ALL-UNNAMED --add-opens=java.base/java.util.concurrent=ALL-UNNAMED --add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED --add-opens=java.base/jdk.internal.ref=ALL-UNNAMED --add-opens=java.base/sun.nio.ch=ALL-UNNAMED --add-opens=java.base/sun.nio.cs=ALL-UNNAMED --add-opens=java.base/sun.security.action=ALL-UNNAMED --add-opens=java.base/sun.util.calendar=ALL-UNNAMED --add-opens=java.security.jgss/sun.security.krb5=ALL-UNNAMED -Djdk.reflect.useDirectMethodHandle=false
- spark.driver.memory: 64g
- spark.dynamicAllocation.enabled: true
- spark.dynamicAllocation.shuffleTracking.enabled: true
- spark.executor.cores: 32
- spark.executor.extraJavaOptions: -Djava.net.preferIPv6Addresses=false -XX:+IgnoreUnrecognizedVMOptions --add-opens=java.base/java.lang=ALL-UNNAMED --add-opens=java.base/java.lang.invoke=ALL-UNNAMED --add-opens=java.base/java.lang.reflect=ALL-UNNAMED --add-opens=java.base/java.io=ALL-UNNAMED --add-opens=java.base/java.net=ALL-UNNAMED --add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/java.util=ALL-UNNAMED --add-opens=java.base/java.util.concurrent=ALL-UNNAMED --add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED --add-opens=java.base/jdk.internal.ref=ALL-UNNAMED --add-opens=java.base/sun.nio.ch=ALL-UNNAMED --add-opens=java.base/sun.nio.cs=ALL-UNNAMED --add-opens=java.base/sun.security.action=ALL-UNNAMED --add-opens=java.base/sun.util.calendar=ALL-UNNAMED --add-opens=java.security.jgss/sun.security.krb5=ALL-UNNAMED -Djdk.reflect.useDirectMethodHandle=false
- spark.executor.id: driver
- spark.executor.memory: 64g
- spark.master: local[*]
- spark.rdd.compress: true
- spark.scheduler.mode: FIFO
- spark.serializer.objectStreamReset: 100
- spark.sql.adaptive.enabled: true
- spark.sql.adaptive.skewJoin.enabled: true
- spark.sql.autoBroadcastJoinThreshold: 1g # 我调整过这个参数,默认0时表现类似
- spark.sql.shuffle.partitions: 64
- spark.submit.deployMode: client
- spark.task.cpus: 32
Hadoop属性
(应为默认值;系统未安装Hadoop)
系统属性
- SPARK_SUBMIT: true
- file.encoding: utf-8
- file.separator: /
- java.class.version: 65.0
- java.runtime.name: OpenJDK Runtime Environment
- java.runtime.version: 21.0.2+13-Ubuntu-122.04.1
- os.arch: amd64
- os.name: Linux
- os.version: 5.15.146.1-microsoft-standard-WSL2
内容的提问来源于stack exchange,提问作者Nathan
相关产品推荐
相关产品推荐

