在Databricks中使用Python加密列并保留原始数据类型
在Databricks中加密DataFrame列并保留原始数据类型的解决方案
要实现DataFrame特定列加密且保留原始数据类型,核心是通过自定义UDF结合对称加密库,针对不同数据类型做序列化/反序列化处理,加密后存储为字符串,解密时再还原为原类型。以下是可落地的实现方案:
实现步骤与代码示例
1. 依赖准备
首先确保安装加密库(如果环境未预装):
%pip install cryptography
2. 加密解密工具实现
使用cryptography的Fernet对称加密方案,针对不同数据类型编写序列化/反序列化逻辑:
from cryptography.fernet import Fernet from pyspark.sql.functions import udf from pyspark.sql.types import ( StringType, IntegerType, FloatType, TimestampType, BooleanType, DateType ) from datetime import datetime, date # 生成密钥(生产环境请用Databricks Secrets存储,禁止硬编码) key = Fernet.generate_key() fernet = Fernet(key) # 加密函数:根据数据类型序列化后加密 def encrypt_val(value, data_type): if value is None: return None # 针对不同类型做序列化处理 if isinstance(data_type, (IntegerType, FloatType, BooleanType)): serialized = str(value).encode() elif isinstance(data_type, TimestampType): serialized = value.strftime("%Y-%m-%d %H:%M:%S.%f").encode() elif isinstance(data_type, DateType): serialized = value.strftime("%Y-%m-%d").encode() else: # StringType及其他类型 serialized = str(value).encode() return fernet.encrypt(serialized).decode() # 解密函数:解密后反序列化为原始类型 def decrypt_val(encrypted_str, data_type): if encrypted_str is None: return None decrypted_bytes = fernet.decrypt(encrypted_str.encode()) decrypted_str = decrypted_bytes.decode() # 反序列化回原始类型 if isinstance(data_type, IntegerType): return int(decrypted_str) elif isinstance(data_type, FloatType): return float(decrypted_str) elif isinstance(data_type, BooleanType): return decrypted_str.lower() == 'true' elif isinstance(data_type, TimestampType): return datetime.strptime(decrypted_str, "%Y-%m-%d %H:%M:%S.%f") elif isinstance(data_type, DateType): return date.fromisoformat(decrypted_str) else: return decrypted_str
3. 注册类型匹配的UDF
为不同数据类型创建对应的加密/解密UDF,确保解密后类型与原始一致:
# 生成加密UDF def get_encrypt_udf(data_type): return udf(lambda x: encrypt_val(x, data_type), StringType()) # 生成解密UDF def get_decrypt_udf(data_type): return udf(lambda x: decrypt_val(x, data_type), data_type)
4. 加密、解密验证
用测试DataFrame验证加密后解密的类型完整性:
# 创建测试DataFrame test_df = spark.createDataFrame([ (1, 123.45, "用户敏感信息", datetime(2024,5,20,10,30), date(2024,5,20), True), (2, 678.90, "机密业务数据", datetime(2024,5,21,14,15), date(2024,5,21), False) ], ["id", "amount", "secret_col", "event_ts", "event_dt", "is_valid"]) # 加密目标列 encrypted_df = test_df \ .withColumn("encrypted_amount", get_encrypt_udf(FloatType())("amount")) \ .withColumn("encrypted_secret", get_encrypt_udf(StringType())("secret_col")) \ .withColumn("encrypted_ts", get_encrypt_udf(TimestampType())("event_ts")) \ .withColumn("encrypted_dt", get_encrypt_udf(DateType())("event_dt")) \ .drop("amount", "secret_col", "event_ts", "event_dt") # 解密验证 decrypted_df = encrypted_df \ .withColumn("decrypted_amount", get_decrypt_udf(FloatType())("encrypted_amount")) \ .withColumn("decrypted_secret", get_decrypt_udf(StringType())("encrypted_secret")) \ .withColumn("decrypted_ts", get_decrypt_udf(TimestampType())("encrypted_ts")) \ .withColumn("decrypted_dt", get_decrypt_udf(DateType())("encrypted_dt")) # 打印Schema验证类型 print("原始Schema:") test_df.printSchema() print("\n解密后Schema:") decrypted_df.printSchema()
关键注意事项
- 密钥安全:生产环境必须用Databricks Secrets管理密钥,示例中硬编码密钥仅用于测试,获取密钥方式:
fernet = Fernet(dbutils.secrets.get("secret_scope", "encryption_key").encode()) - 性能优化:处理超大数据集时,改用
pandas_udf替代普通UDF,提升并行处理效率 - 类型扩展:如果有自定义数据类型,只需在加密/解密函数中添加对应的序列化/反序列化逻辑即可
- 空值兼容:代码中已处理空值场景,避免加密空值引发异常
内容的提问来源于stack exchange,提问作者practicalGuy
相关产品推荐
相关产品推荐

