Databricks使用sqlContext遍历列表生成新列时unionAll结果被覆盖
问题根因
你遇到的最后一条结果覆盖所有条目的问题,主要由3个错误导致:
- SQL语法错误:你的查询语句缺少
WHERE关键字,LIKE条件没有指定匹配的字段,FROM config.uploader LIKE LOWER(...)会被Spark SQL解析为给表起别名为LIKE,传入的列表参数f根本没有参与过滤,每次查询实际返回的逻辑不符合预期。 - 循环逻辑错误:第一次循环时
df变量没有被正确初始化就调用unionAll,且你写的if/else分支逻辑完全重复,没有实际意义。 - 实现方式低效:循环发起多次SQL查询再做union,会产生大量冗余IO和计算,还容易触发Spark懒执行时的变量绑定异常。
最优实现方案
完全不需要循环遍历列表发起多次查询,直接在SQL层拼接所有匹配条件,一次查询即可拿到所有目标结果,性能远高于循环union,也不会出现结果覆盖问题:
# 拼接所有LIKE匹配条件 like_clause = " OR ".join([f"LOWER(tablename) LIKE LOWER('%{keyword}%')" for keyword in row_list]) query_sql = f""" SELECT concat(domain, '_', quality, '.', sink) AS datalake_table, tablename FROM config.uploader WHERE {like_clause} """ result_df = sqlContext.sql(query_sql)
如果你的匹配规则是精确匹配表名,还可以把LIKE换成IN条件,查询效率更高:
# 精确匹配场景用IN in_clause = ",".join([f"'{keyword}'" for keyword in row_list]) query_sql = f""" SELECT concat(domain, '_', quality, '.', sink) AS datalake_table, tablename FROM config.uploader WHERE LOWER(tablename) IN ({in_clause}) """ result_df = sqlContext.sql(query_sql)
循环方式修复方案(不推荐)
如果坚持要用循环union的逻辑,需要修正初始化逻辑和SQL语法,代码如下:
result_df = None for keyword in row_list: tmp_df = sqlContext.sql(f""" SELECT concat(domain, '_', quality, '.', sink) AS datalake_table, tablename FROM config.uploader WHERE LOWER(tablename) LIKE LOWER('%{keyword}%') """) if result_df is None: result_df = tmp_df else: result_df = result_df.unionAll(tmp_df)
注意:Spark的unionAll/union默认按列位置合并、不会去重,如果需要去重可以在合并完成后调用.distinct()方法。
内容的提问来源于stack exchange,提问作者Shawn N
相关产品推荐
相关产品推荐

