如何使用PySpark并行执行多表INSERT语句 无需使用多线程
PySpark 无多线程实现多INSERT语句并行执行方案
默认Spark采用FIFO(先进先出)调度模式,你当前的写法每次调用spark.sql()执行INSERT操作时会立即触发作业提交,前一个作业执行完成后才会启动下一个,因此是串行执行。你可以利用Spark原生的公平调度能力实现并行,无需自己编写多线程逻辑,具体方案如下:
1. 调整Spark调度模式为公平调度
初始化SparkSession时添加公平调度配置,允许Spark同时调度多个独立作业:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("parallel_insert_task") \ # 开启公平调度模式 .config("spark.scheduler.mode", "FAIR") \ .getOrCreate()
2. 预构建所有待写入的数据集
先把每个表对应的查询逻辑提取出来生成DataFrame,这一步是Spark的转换操作,属于惰性求值,只会生成逻辑执行计划,不会触发实际计算:
# 统一维护所有表的同步配置,方便后续扩展 table_sync_config = [ { "staging_table": "tbl1", "target_table": "Cls.tbl1", "select_columns": ["Contract", "Name"] }, { "staging_table": "tbl2", "target_table": "Cls.tbl2", "select_columns": ["Contract", "Name"] } ] # 预生成所有待写入的DataFrame wait_write_df = {} for config in table_sync_config: stg_tbl = config["staging_table"] target_tbl = config["target_table"] query_sql = f""" SELECT s.Contract, s.Name FROM {stg_tbl} AS s LEFT JOIN {target_tbl} AS c ON s.Contract = c.Contract AND s.Adj = c.Adj WHERE c.Contract IS NULL """ wait_write_df[target_tbl] = spark.sql(query_sql)
3. 批量触发写入操作
循环执行写入操作时,每个写入会提交一个独立作业,公平调度器会根据集群空闲资源自动并行调度多个作业执行:
for target_table, df in wait_write_df.items(): df.write.mode("append").insertInto(target_table)
注意事项
- 并行度由集群可用资源决定,只要集群有空闲CPU核心,多个写入作业就会同时执行
- 如果多个表的查询逻辑有公共上游依赖,可以提前对公共依赖的DataFrame执行
cache()操作,避免重复计算 - 可以通过
spark.scheduler.allocation.file配置自定义调度池,给不同优先级的表分配不同的资源权重
内容的提问来源于stack exchange,提问作者Rocky1989
相关产品推荐
相关产品推荐

