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

PySpark如何加密JSON struct中pii字段并写回以保留原JSON结构

解决方案

你无需将JSON打平再重构,PySpark针对struct类型字段提供了直接更新嵌套字段的能力,两种可行方案如下:

方案1:使用withField直接更新嵌套字段(适用于Spark 3.1及以上版本)

该方案为最优解,可直接对struct类型的json列的pii字段做修改,其他字段完全保留原值,不会改动原JSON结构:

from pyspark.sql import functions as F
from pyspark.sql.types import StringType

# 定义加密UDF,替换为你的实际加密逻辑
@F.udf(returnType=StringType())
def encrypt_pii(raw_pii: str) -> str:
    # 示例加密逻辑
    return f"encrypted_{raw_pii}"

# 直接更新json结构体中的pii字段
df_updated = df.withColumn(
    "json",
    F.col("json").withField("pii", encrypt_pii(F.col("json.pii")))
)

方案2:手动重构struct(适用于Spark 3.0及以下低版本)

如果你的Spark版本不支持withField方法,可以手动组装struct结构体,仅修改pii字段,其余字段原样保留:

from pyspark.sql import functions as F
from pyspark.sql.types import StringType

@F.udf(returnType=StringType())
def encrypt_pii(raw_pii: str) -> str:
    return f"encrypted_{raw_pii}"

# 自动获取json结构体的所有字段名,避免手动遗漏
json_fields = df.schema["json"].dataType.fieldNames()
struct_fields = []
for field in json_fields:
    if field == "pii":
        struct_fields.append(encrypt_pii(F.col(f"json.{field}")).alias(field))
    else:
        struct_fields.append(F.col(f"json.{field}").alias(field))

# 重新构建json struct列
df_updated = df.withColumn("json", F.struct(*struct_fields))

修改完成后可通过df_updated.select(F.to_json("json")).show(truncate=False)验证最终JSON结构,确认仅pii字段被加密,其余结构与原值完全一致。

内容的提问来源于stack exchange,提问作者T.UK

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 18:33:01