如何在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
相关产品推荐
相关产品推荐

