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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 15:07:29