如何在PySpark函数中引入JSON字符串变量并实现数据过滤
Databricks中PySpark基于Excel动态JSON规则的数据过滤与GSBER校验方案
全流程思路
先把Excel里的JSON规则读出来解析成Python字典(比直接转DataFrame更灵活,毕竟是配置规则),然后根据规则生成过滤逻辑(两种实现方式:DataFrame关联过滤/动态SQL),最后加上GSBER的历史表校验。
先修正你的JSON规则(原示例有语法错误)
以下是修正后的合法JSON,后续代码基于此:
// 用逻辑表过滤数据表 // - 规则键由类型和列名按顺序定义,所有条件AND组合 // key[0]:匹配数据表的"GSBER_BUSINESS_AREA"与逻辑表的"GSBER" // key[1]:判断数据表的"GJAHR_FISCAL_YEAR" >= 逻辑表的"FY range start" // AND 数据表的"GJAHR_FISCAL_YEAR" <= 逻辑表的"FY range end" // - FilterAction:"include"表示保留匹配数据,"EXCLUDE"表示排除匹配数据 { "process": "Filter", "logicTableNm": "WC/Lab-ERP", "logicKeyTp": [ {"match": ["GSBER", "GSBER_BUSINESS_AREA"]}, {"inRange": ["GJAHR_FISCAL_YEAR", "FY range start", "FY range end"]} ], "FilterAction": "include" }
步骤1:从Excel读取并解析JSON规则
在Databricks中可以用pandas读取Excel,再解析JSON:
import pandas as pd import json # 假设Excel文件已挂载到DBFS,替换为你的实际路径 excel_path = "/dbfs/mnt/your-storage/filter_rules.xlsx" # 读取Excel中存储规则的sheet(假设sheet名为FilterRules) df_excel = pd.read_excel(excel_path, sheet_name="FilterRules") # 提取第一行的JSON规则(假设规则存在于RuleJSON列) rule_json_str = df_excel.loc[0, "RuleJSON"] # 解析为Python字典 filter_rule = json.loads(rule_json_str)
步骤2:核心问题实现——两种过滤方案
方案一:PySpark DataFrame关联+过滤(推荐,类型安全)
直接通过DataFrame的join和过滤实现规则,避免SQL拼接的风险:
from pyspark.sql.functions import col # 加载目标业务数据表(替换为你的实际表名) target_df = spark.table("erp.business_data") # 加载JSON中指定的逻辑表 logic_df = spark.table(filter_rule["logicTableNm"]) # 构建关联条件列表 join_conditions = [] for rule in filter_rule["logicKeyTp"]: if "match" in rule: # match规则:逻辑表列 -> 目标表列 logic_col, target_col = rule["match"] join_conditions.append(col(f"l.{logic_col}") == col(f"t.{target_col}")) elif "inRange" in rule: # inRange规则:目标表列 >= 逻辑表起始列 AND 目标表列 <= 逻辑表结束列 target_col, range_start, range_end = rule["inRange"] join_conditions.append(col(f"t.{target_col}") >= col(f"l.{range_start}")) join_conditions.append(col(f"t.{target_col}") <= col(f"l.{range_end}")) # 关联目标表与逻辑表 joined_df = target_df.alias("t").join(logic_df.alias("l"), on=join_conditions, how="inner") # 根据FilterAction处理结果 action = filter_rule.get("FilterAction", "include").lower() if action == "exclude": # 排除匹配到的记录,用left_anti关联 filtered_df = target_df.alias("t").join(joined_df.alias("j"), on=target_df["GSBER_BUSINESS_AREA"] == joined_df["GSBER_BUSINESS_AREA"], how="left_anti") else: # 保留匹配记录,去重避免逻辑表多条数据导致重复 filtered_df = joined_df.select("t.*").distinct() # 查看结果 filtered_df.display()
方案二:动态生成Spark SQL WHERE子句(适合熟悉SQL的场景)
解析规则生成SQL条件字符串,直接执行SQL查询:
# 构建WHERE条件 where_clauses = [] logic_table = filter_rule["logicTableNm"] for rule in filter_rule["logicKeyTp"]: if "match" in rule: logic_col, target_col = rule["match"] where_clauses.append(f"t.{target_col} = l.{logic_col}") elif "inRange" in rule: target_col, range_start, range_end = rule["inRange"] where_clauses.append(f"t.{target_col} >= l.{range_start}") where_clauses.append(f"t.{target_col} <= l.{range_end}") where_str = " AND ".join(where_clauses) action = filter_rule.get("FilterAction", "include").lower() # 生成完整SQL if action == "include": sql_query = f""" SELECT DISTINCT t.* FROM erp.business_data t JOIN (SELECT * FROM {logic_table}) l ON {where_str} """ else: sql_query = f""" SELECT t.* FROM erp.business_data t LEFT JOIN (SELECT * FROM {logic_table}) l ON {where_str} WHERE l.GSBER IS NULL """ # 执行SQL filtered_df = spark.sql(sql_query) filtered_df.display()
步骤3:GSBER字段历史表校验
在过滤后的数据基础上,关联历史逻辑表判断GSBER是否存在:
# 加载历史逻辑表(替换为你的实际表名) history_logic_df = spark.table("erp.history_logic_table") # 左关联历史表,标记是否存在 validated_df = filtered_df.alias("t").join( history_logic_df.alias("h"), on=col("t.GSBER_BUSINESS_AREA") == col("h.GSBER"), how="left" ).withColumn( "GSBER_IS_EXISTS_IN_HISTORY", col("h.GSBER").isNotNull() ).select("t.*", "GSBER_IS_EXISTS_IN_HISTORY") # 查看校验结果 validated_df.display()
注意事项
- 维护Excel中的JSON时,务必保证语法正确,可在线JSON校验工具验证
- 确保逻辑表与目标表的关联字段类型一致,避免关联失败
- 动态生成SQL时,禁止外部用户修改规则,防止SQL注入风险
- Databricks中读取云存储(S3/ADLS)的Excel文件,需先配置存储挂载或直接使用云路径
内容的提问来源于stack exchange,提问作者Bruce Jenks
相关产品推荐
相关产品推荐

