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

如何将Azure Key Vault获取的AES密钥用于Spark aes_encrypt加密DataFrame列?

从Azure Key Vault获取AES密钥用于Spark aes_encrypt函数

可以实现,但需要注意两个核心要点:Spark的aes_encrypt函数要求传入原始AES密钥字节的Base64编码字符串,而从Key Vault获取的密钥对象无法直接使用,需先提取原始密钥材料(仅当密钥设置为可导出时才能获取)。

方法一:使用可导出的AES密钥(推荐,效率更高)

如果在Azure Key Vault创建AES密钥时设置了exportable=True,可直接提取原始密钥字节,转换为Base64字符串后传入Spark函数。

代码实现

from pyspark.sql import functions as F
from azure.identity import DefaultAzureCredential
from azure.keyvault.keys import KeyClient
import base64

# 1. 从Key Vault获取可导出的AES密钥
credential = DefaultAzureCredential()
key_client = KeyClient(vault_url="https://my-key-vault.vault.azure.net/", credential=credential)
# 需确保密钥创建时设置了exportable=True
key = key_client.get_key("MyKey")

# 2. 提取原始密钥字节并转为Base64编码字符串
key_base64 = base64.b64encode(key.key).decode("utf-8")

# 3. 用Spark aes_encrypt加密DataFrame列
# 数值类型需先转为字符串/二进制类型
parDF1 = parDF1.withColumn(
    'encrypted_value', 
    F.aes_encrypt(F.col("tripDistance").cast("string"), key_base64, "GCM")
)

方法二:密钥不可导出时的处理(效率较低)

如果Key Vault中的AES密钥设置为不可导出(exportable=False),无法直接获取原始密钥字节,只能通过Azure Key Vault的加密API逐行加密数据。此时需自定义Spark UDF,调用CryptographyClient完成加密。

代码实现

from pyspark.sql import functions as F
from pyspark.sql.types import BinaryType
from azure.identity import DefaultAzureCredential
from azure.keyvault.keys import KeyClient
from azure.keyvault.keys.crypto import CryptographyClient, EncryptionAlgorithm

# 1. 初始化Key Vault客户端和加密客户端
credential = DefaultAzureCredential()
key_client = KeyClient(vault_url="https://my-key-vault.vault.azure.net/", credential=credential)
key = key_client.get_key("MyKey")
crypto_client = CryptographyClient(key, credential=credential)

# 2. 定义自定义加密UDF
def encrypt_data(data):
    if not data:
        return None
    # 将数据转为字节(数值类型先转字符串再编码)
    plaintext = str(data).encode("utf-8")
    # 调用Key Vault加密接口
    result = crypto_client.encrypt(EncryptionAlgorithm.a128_gcm, plaintext)
    # 返回包含密文、IV和标签的组合数据(解密时需完整传入)
    return result.ciphertext + result.iv + result.tag

encrypt_udf = F.udf(encrypt_data, BinaryType())

# 3. 应用UDF加密列
parDF1 = parDF1.withColumn('encrypted_value', encrypt_udf(F.col("tripDistance")))

注意事项

  • 方法二的UDF在分布式环境运行时,每个Executor都会创建CryptographyClient实例,需确保Executor具备访问Key Vault的权限(通过MSI或服务主体)。
  • 解密时需对应使用crypto_client.decrypt方法,传入加密时的算法以及密文+IV+标签的组合数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 03:54:29