减少Spark作业启动次数的最佳实践:使用union方案是否可行?
结论
合并为单个Union作业的方案是不是最佳实践没有统一答案,绝大多数场景下值得优先考虑,但需要结合业务场景、数据量和执行计划表现判断,具体分析如下:
优先选Union方案的核心原因
- 作业启动开销的节省远大于单作业复杂度提升带来的额外消耗:Spark每次提交作业都会产生调度器排队、Driver与Executor通信、资源申请(开启动态资源分配时)等固定开销,单次开销看起来不大,但循环次数超过5次之后累加的开销会远高于单作业多算的时间。比如循环100次每次作业跑1s,启动开销0.5s的话,总耗时会达到150s,换成Union方案哪怕单作业跑30s,性能也提升了4倍。
- Spark优化器对Union逻辑的适配已经非常成熟:Catalyst优化器会自动对合并后的执行计划做谓词下推、列裁剪、分区修剪等通用优化,多数情况下合并后的总资源消耗反而低于多个小作业的总和,不会出现复杂度提升导致性能暴跌的问题。
- 稳定性更优:多次提交作业意味着多次和集群调度器交互,中间任意一次作业出现资源不足、节点故障等问题都会导致整个流程失败,合并为单作业的容错、重试成本都更低。
不适合用Union方案的边界场景
遇到以下情况时,反而建议拆分作业运行:
- 合并后的数据量远超集群承载上限:如果Union的子查询数量过多、单条子查询数据量达到TB级,合并后单作业的Shuffle数据量太大,很容易出现Executor OOM、Shuffle写失败等问题,这种场景下拆分为多个小作业分批跑稳定性更好。
- 执行计划出现明显退化:如果每个子查询的逻辑差异极大,合并后优化器出现判断失误,比如原本部分子查询可以走索引、分区裁剪,合并后反而走了全表扫描,这种情况可以先对比两种方案的
explain执行计划,确认没有退化再使用Union方案。 - 有中间结果复用需求:如果循环生成的部分子查询结果后续流程还要复用,那拆分作业后缓存中间结果的收益会高于合并作业的收益。
实操建议
如果循环次数低于3次,两种方案的性能差异极小,优先选代码可读性更高的实现即可;如果循环次数更多,先拿小批量数据验证两种方案的耗时和稳定性,再做选择。
附两种实现的代码参考:
循环提交作业实现
def test_loop(spark): all_datas = [] for i in ['CZ12905K01', 'CZ12809WRH', 'CZ129086RP']: all_datas.extend(spark.sql(f""" select * from data where id=='{i}' """).collect()) # 每次循环启动一个作业 return all_datas
Union合并单作业实现
def test_union(spark): full_request = None for i in ['CZ12905K01', 'CZ12809WRH' ,'CZ129086RP']: q = f""" select '{i}' ID,* from data where leh_be_lot_id=='{i}' """ partial_df = spark.sql(q) if not full_request: full_request = partial_df else: full_request = full_request.union(partial_df) return full_request.collect() # 仅启动一个作业
注:上述为简化示例(实际可直接使用in语法实现查询),真实业务场景的查询逻辑会更加复杂。
内容的提问来源于stack exchange,提问作者Philippe Prados
相关产品推荐
相关产品推荐

