You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 04:57:33