在Databricks PySpark中执行动态生成的Union字符串语句
在Databricks PySpark中执行动态构建的Union字符串语句
问题说明
我在Databricks笔记本里用PySpark,已经把Union语句动态构建成了字符串变量,比如输出的df533.union(df534).union(df535).union(df536),但卡壳在执行这串代码的环节。
以下是我构建这个字符串的代码:
from pyspark.sql.functions import col, concat, lit, concat_ws, overlay df1 = df.filter((col("vchDataSection") == "AccountMasterInfo") & (col("bActive") == 1)).withColumn("dfs", concat(lit(".union(df"), col("iRuleid"), lit(")"))) df2 = df1.agg(concat_ws("", collect_list(col("dfs")))).withColumnRenamed("concat_ws(, collect_list(dfs))", "AccInfoRules").withColumn("replacestr", lit("")) df3 = df2.select(overlay("AccInfoRules", "replacestr", 1, 7).alias("overlayed")) var_a = df3.collect() var_a = var_a[0].__getitem__('overlayed') var_b = var_a.replace(')', '', 1) print(var_b)
执行后输出:
df533.union(df534).union(df535).union(df536)
解决办法
方式1:直接用eval()执行字符串
如果能保证字符串里的内容绝对安全(没有恶意代码),可以直接用eval()运行这个字符串得到合并后的DataFrame:
result_df = eval(var_b) # 查看合并结果 result_df.show()
⚠️ 注意:eval()会执行字符串里的任意代码,要是字符串来源不可信,会有安全风险,谨慎使用。
方式2:更安全的动态合并(推荐)
不想用eval()的话,可以提取需要合并的DataFrame对象,再循环合并:
- 先从字符串里提取所有要合并的DataFrame名称:
import re # 匹配所有df开头带数字的名称 df_names = re.findall(r'df\d+', var_b) # 从全局变量里获取对应的DataFrame对象 dfs_to_union = [globals()[df_name] for df_name in df_names]
- 用循环或者
reduce来合并所有DataFrame:
# 方法A:循环合并 result_df = dfs_to_union[0] for single_df in dfs_to_union[1:]: result_df = result_df.union(single_df) # 方法B:用reduce更简洁 from functools import reduce result_df = reduce(lambda df_a, df_b: df_a.union(df_b), dfs_to_union)
这种方式完全不需要执行字符串代码,安全性更高,也更符合PySpark的使用习惯。
内容的提问来源于stack exchange,提问作者user2225726
相关产品推荐
相关产品推荐

