Spark中如何高效判断一个数组是否为另一个数组的子集
Spark计算场景下数组子集判断实现方案
在Spark中判断数组A是否为数组B的子集(即A中所有元素都存在于B中),需要兼容NULL值、重复元素、空输入等边界场景,以下基于提供的测试DataFrame给出可落地的实现方案。
测试数据
测试DataFrame变量名为df,覆盖了常规匹配、不匹配、NULL元素、重复元素、全NULL输入等场景,表结构如下:
+------------+------------+ |look_in |look_for | +------------+------------+ |[a, b, c] |[a] | |[a, b, c] |[d] | |[a, b, c] |[a, b] | |[a, b, c] |[c, d] | |[a, b, c] |[a, b, c] | |[a, b, c] |[a, NULL] | |[a, b, NULL]|[a, NULL] | |[a, b, NULL]|[a] | |[a, b, NULL]|[NULL] | |[a, b, c] |NULL | |NULL |[a] | |NULL |NULL | |[a, b, c] |[a, a] | |[a, a, a] |[a] | |[a, a, a] |[a, a, a] | |[a, a, a] |[a, a, NULL]| |[a, a, NULL]|[a, a, a] | |[a, a, NULL]|[a, a, NULL]| +------------+------------+
DataFrame构造代码:
from pyspark.sql import functions as F df = spark.createDataFrame( [(['a', 'b', 'c'], ['a']), (['a', 'b', 'c'], ['d']), (['a', 'b', 'c'], ['a', 'b']), (['a', 'b', 'c'], ['c', 'd']), (['a', 'b', 'c'], ['a', 'b', 'c']), (['a', 'b', 'c'], ['a', None]), (['a', 'b',None], ['a', None]), (['a', 'b',None], ['a']), (['a', 'b',None], [None]), (['a', 'b', 'c'], None), (None, ['a']), (None, None), (['a', 'b', 'c'], ['a', 'a']), (['a', 'a', 'a'], ['a']), (['a', 'a', 'a'], ['a', 'a', 'a']), (['a', 'a', 'a'], ['a', 'a',None]), (['a', 'a',None], ['a', 'a', 'a']), (['a', 'a',None], ['a', 'a',None])], ['look_in', 'look_for'])
实现方案
方案1:内置集合函数判断(性能最优,Spark 2.4+支持)
核心逻辑:对数组去重后计算交集,通过交集长度和待查数组去重长度的对比判断子集关系,单独兼容NULL元素匹配逻辑。
res1 = df.withColumn( "is_subset", # 输入数组本身为NULL时返回NULL F.when(F.col("look_in").isNull() | F.col("look_for").isNull(), None) .when( # 校验待查数组中的NULL元素是否在父数组存在 F.expr("exists(look_for, x -> x is null)") & ~F.expr("exists(look_in, x -> x is null)"), False ) .otherwise( # 过滤NULL后去重,判断交集长度 F.size(F.array_intersect(F.expr("filter(look_in, x -> x is not null)"), F.expr("filter(look_for, x -> x is not null)"))) == F.size(F.expr("array_distinct(filter(look_for, x -> x is not null))")) ) )
- 优点:全部使用Spark原生内置函数,无UDF序列化开销,Catalyst引擎可全程优化,大数据量下性能最高,兼容Spark 2.4及以上版本
- 缺点:NULL处理逻辑相对繁琐
方案2:高阶函数forall遍历判断(写法最简洁,Spark 2.4+支持)
核心逻辑:用forall高阶函数遍历待查数组的每个元素,判断元素是否存在于父数组中,通过coalesce处理NULL值的三值逻辑问题。
res2 = df.withColumn( "is_subset", F.when(F.col("look_in").isNull() | F.col("look_for").isNull(), None) .otherwise(F.expr("forall(look_for, x -> coalesce(array_contains(look_in, x), false))")) )
- 优点:代码极简,逻辑直观,自动兼容重复元素场景,不需要额外做去重处理
- 缺点:时间复杂度为O(n*m)(n、m分别为两个数组的长度),数组长度极大时性能略低于方案1,日常亿级数据量下无明显瓶颈
方案3:Python UDF集合判断(兼容低版本Spark)
核心逻辑:自定义Python函数将数组转为集合,利用Python集合的issubset方法做判断,使用前需导入BooleanType类型。
from pyspark.sql.types import BooleanType @F.udf(returnType=BooleanType()) def array_subset(look_in, look_for): if look_in is None or look_for is None: return None # 单独处理NULL元素匹配 if None in look_for and None not in look_in: return False # 过滤NULL后转集合判断 in_set = {i for i in look_in if i is not None} for_set = {i for i in look_for if i is not None} return for_set.issubset(in_set) res3 = df.withColumn("is_subset", array_subset("look_in", "look_for"))
- 优点:逻辑灵活可自定义,兼容所有Spark版本
- 缺点:Python UDF存在跨进程序列化/反序列化开销,性能比原生函数低3~10倍,TB级超大数据量下不推荐使用
预期运行结果
三个方案返回的判断结果完全一致,具体如下:
+------------+------------+---------+ |look_in |look_for |is_subset| +------------+------------+---------+ |[a, b, c] |[a] |true | |[a, b, c] |[d] |false | |[a, b, c] |[a, b] |true | |[a, b, c] |[c, d] |false | |[a, b, c] |[a, b, c] |true | |[a, b, c] |[a, NULL] |false | |[a, b, NULL]|[a, NULL] |true | |[a, b, NULL]|[a] |true | |[a, b, NULL]|[NULL] |true | |[a, b, c] |NULL |NULL | |NULL |[a] |NULL | |NULL |NULL |NULL | |[a, b, c] |[a, a] |true | |[a, a, a] |[a] |true | |[a, a, a] |[a, a, a] |true | |[a, a, a] |[a, a, NULL]|false | |[a, a, NULL]|[a, a, a] |true | |[a, a, NULL]|[a, a, NULL]|true | +------------+------------+---------+
内容的提问来源于stack exchange,提问作者ZygD
相关产品推荐
相关产品推荐

