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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 15:45:25