Spark Scala如何在withColumn子句中执行存储于列的SQL查询
错误原因
spark.sql()方法仅支持传入字符串格式的SQL语句,你传入的col("QUERY")是Spark的Column类型对象,两者类型不匹配,因此直接抛出类型错误。
能否用UDF实现
不能直接通过UDF执行spark.sql()完成需求。UDF运行在分布式Executor节点上,而spark.sql()依赖的SparkSession仅在Driver端有效,Executor端无法正常调用该对象,强制使用会出现序列化错误、空指针异常等问题。
正确实现方案
由于你的查询都是单字段limit 1的轻量查询,可以先在Driver端批量执行所有查询生成结果映射,再把结果关联回原DataFrame,实现代码如下:
from pyspark.sql.functions import udf, col from pyspark.sql.types import StringType # 1. 从原DF中提取字段名和对应查询,收集到Driver端 query_mapping = DF.select("DIFFCOLUMNNAME", "QUERY").collect() # 2. 逐条执行查询,生成 字段名:查询结果 的映射字典 result_dict = {} for item in query_mapping: col_name = item["DIFFCOLUMNNAME"] exec_sql = item["QUERY"] # 执行SQL取返回的第一个值 exec_result = spark.sql(exec_sql).collect()[0][0] result_dict[col_name] = exec_result # 3. 构建映射UDF,把结果写入原DF @udf(returnType=StringType()) def map_result(col_name): return str(result_dict.get(col_name)) final_df = DF.withColumn("QUERYRESULT", map_result(col("DIFFCOLUMNNAME"))) # 查看结果 final_df.show()
注意事项
该方案适用于QUERY列行数不多的场景,如果你的原DF行数过万,可考虑批量拼接SQL减少查询次数优化效率。
内容的提问来源于stack exchange,提问作者pinksrider
相关产品推荐
相关产品推荐

