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

