Spark中Union与Stack的性能对比及实际数据场景咨询
Spark中Union与Stack的性能对比及实际数据场景咨询
嘿,我来帮你拆解Spark里Union和Stack在你这个实体关联场景下的用法和性能差异~咱先从你的数据结构说起,你给出的DataFrame是这样的:
df = spark.createDataFrame( [ ['A', 'A06', 'B', 'B02', '202412'], ['A', 'A04', 'B', 'B03', '202501'], ['B', 'B01', 'C', 'C02', '202411'], ['B', 'B03', 'A', 'A06', '202502'] ], 'entity_code: string, entity_rollup: string, target_entity_code: string, target_entity_rollup: string, period: string' ) df.createOrReplaceTempView('v1')
当你需要把实体和目标实体的双向关系展开时,就会纠结用Union还是Stack,对吧?下面咱分别唠唠:
一、Union的用法与性能特点
Union的核心是把两个结构完全一致的DataFrame拼接在一起。比如你要得到双向关系,得先把原数据的实体和目标实体列交换,生成一个新的DataFrame,再和原df做Union:
# 生成交换列的DataFrame reversed_df = df.select( col("target_entity_code").alias("entity_code"), col("target_entity_rollup").alias("entity_rollup"), col("entity_code").alias("target_entity_code"), col("entity_rollup").alias("target_entity_rollup"), col("period") ) # 合并原df和反转后的df union_result = df.union(reversed_df)
性能上要注意这几点:
- Union属于宽依赖,Spark会保留两个输入DataFrame的分区结构,不会自动合并分区(除非你手动加
coalesce或repartition) - 它需要扫描两遍数据:一遍处理原df,一遍处理反转后的reversed_df,数据量越大,重复扫描的开销就越明显
- 适合数据量较小、或者需要保留原数据分区逻辑的场景,比如小批量的实体关系合并
二、Stack的用法与性能特点
Stack是Spark SQL里的内置函数,核心是在单遍扫描数据时,把一条记录拆成多行,完全不需要额外生成一个反转的DataFrame。针对你的场景,用Stack的SQL写法是这样的:
SELECT period, stack(2, entity_code, entity_rollup, target_entity_code, target_entity_rollup, target_entity_code, target_entity_rollup, entity_code, entity_rollup ) AS (new_entity_code, new_entity_rollup, new_target_code, new_target_rollup) FROM v1
性能优势就很突出了:
- Stack是窄依赖,Spark可以在同一个分区内直接处理每条记录,把一行拆成两行,全程不需要额外的shuffle或者重复扫描数据
- 大数据量下,单遍扫描的效率比Union的两遍扫描高很多,能节省不少IO和计算资源
- 写法更简洁,不需要额外生成中间DataFrame,代码维护起来也更省心
三、你的实体关联场景该选哪个?
结合你这个实体-目标实体的周期数据场景:
- 如果你的数据量不大(比如百万级以内),Union和Stack的性能差异感知不明显,选哪个看你代码习惯就行
- 但如果是千万级甚至更大的数据量,优先选Stack,因为它避免了重复计算和扫描,能显著降低作业的运行时间
- 另外,如果你需要保留原DataFrame的所有列结构(包括顺序),Union可以直接对齐结构,而Stack需要你手动指定所有要展开的列,这点要根据你的需求来权衡
备注:内容来源于stack exchange,提问作者Dhruv
相关产品推荐
相关产品推荐

