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

如何在Azure Synapse PySpark Notebook中加密Parquet特定列数据?

在Azure Synapse PySpark中写入Parquet时加密PII列的实现方案

一、自定义UDF实现轻量加密(适合小型场景)

这种方式无需依赖外部密钥管理服务,使用对称加密算法快速处理PII列,适合对加密复杂度要求不高的场景。

实现步骤

  • 导入加密依赖库,定义加密逻辑并封装为Spark UDF
  • 对目标PII列应用加密UDF,替换或保留明文列
  • 将加密后的DataFrame写入Parquet文件
# 导入必要库
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType
from cryptography.fernet import Fernet

# 生成加密密钥(示例用,实际需从安全存储获取,禁止硬编码)
encryption_key = Fernet.generate_key()
fernet_handler = Fernet(encryption_key)

# 定义加密函数
def encrypt_pii_field(data):
    if data is None:
        return None
    return fernet_handler.encrypt(data.encode()).decode()

# 注册Spark UDF
encrypt_pii_udf = udf(encrypt_pii_field, StringType())

# 加载原始数据(示例为CSV,可替换为你的数据源)
raw_df = spark.read.csv("abfss://container@storageaccount.dfs.core.windows.net/raw_data.csv", header=True)

# 加密指定PII列,比如邮箱、手机号
encrypted_df = raw_df.withColumn("encrypted_email", encrypt_pii_udf(raw_df["email"])) \
                     .withColumn("encrypted_phone", encrypt_pii_udf(raw_df["phone"])) \
                     .drop("email", "phone")  # 可选:删除明文列避免泄露

# 写入加密后的Parquet文件
encrypted_df.write.mode("overwrite").parquet("abfss://container@storageaccount.dfs.core.windows.net/encrypted_parquet/")

关键注意点

  • 密钥必须妥善存储,禁止硬编码在代码中,可后续迁移到Azure Key Vault管理
  • 如需解密,可反向实现解密UDF,使用同一密钥还原数据

二、结合Azure Key Vault实现企业级加密(合规场景)

这种方式符合企业安全规范,密钥由Azure Key Vault统一托管,避免密钥泄露风险,适合对数据合规性要求高的场景。

实现步骤

  • 在Azure Key Vault中创建并存储对称加密密钥
  • 配置Synapse与Key Vault的链接服务,赋予Synapse访问密钥的权限
  • 通过mssparkutils从Key Vault获取密钥,实现加密逻辑
  • 处理PII列并写入Parquet
# 导入依赖库
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType
from cryptography.fernet import Fernet
import mssparkutils

# 从Azure Key Vault获取加密密钥(需提前配置链接服务)
key_vault_name = "your-key-vault-name"
secret_name = "pii-encryption-key"
# 获取密钥(自动处理权限验证)
encryption_key = mssparkutils.credentials.getSecret(key_vault_name, secret_name)
fernet_handler = Fernet(encryption_key.encode())

# 定义加密UDF
def encrypt_pii_field(data):
    if data is None:
        return None
    return fernet_handler.encrypt(data.encode()).decode()

encrypt_pii_udf = udf(encrypt_pii_field, StringType())

# 加载数据并加密PII列(示例从Synapse表读取)
raw_df = spark.read.table("raw_database.raw_data_table")
encrypted_df = raw_df.withColumn("encrypted_ssn", encrypt_pii_udf(raw_df["ssn"])) \
                     .drop("ssn")

# 写入Parquet文件
encrypted_df.write.mode("append").parquet("abfss://container@storageaccount.dfs.core.windows.net/enterprise_encrypted_parquet/")

额外优化建议

  • 给Synapse工作区的服务主体配置Key Vault的密钥读取权限,避免权限报错
  • 可结合Azure存储账户的静态服务端加密,实现“列级加密+存储加密”的双重保障
  • 如需非对称加密,可使用azure-keyvault-keys库的RsaKeyCryptographyClient实现

内容的提问来源于stack exchange,提问作者Prayag15

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 15:01:20