AWS Glue运行Map.apply解密列时触发PicklingError问题咨询
问题背景
在AWS Glue作业中调用Map.apply对DynamicFrame指定列执行自定义解密逻辑时,抛出如下错误:
PicklingError: Could not serialize object: TypeError: can't pickle _ModuleWithDeprecations objects
相同依赖版本下,同一段代码在本地环境可正常运行,初步判断问题与Spark、Glue的脚本打包分发机制,以及基于cryptography库实现的解密UDF有关,需确认是否存在可直接实现需求的配置或写法,还是该问题属于AWS Glue Python版本固有局限,必须切换为预打包依赖的Jar包、使用Scala编写逻辑才能适配。
复现代码
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from base64 import b64decode, b64encode from cryptography.fernet import Fernet from cryptography.hazmat.primitives import hashes from cryptography.hazmat.primitives.kdf.pbkdf2 import PBKDF2HMAC from cryptography.hazmat.backends import default_backend from cryptography.hazmat.primitives.padding import PKCS7 from cryptography.hazmat.primitives.ciphers import Cipher, algorithms, modes KEY = b'our-secret-key-value' def decrypt_pbe_with_hmac_sha512_aes_256(obj: str) -> str: # re-generate key from encrypted_obj = b64decode(obj) salt = encrypted_obj[0:16] iv = encrypted_obj[16:32] cypher_text = encrypted_obj[32:] kdf = PBKDF2HMAC(hashes.SHA512(), 32, salt, 1000, backend=default_backend()) key = kdf.derive(KEY) # decrypt cipher = Cipher(algorithms.AES(key), modes.CBC(iv), backend=default_backend()) decryptor = cipher.decryptor() padded_text = decryptor.update(cypher_text) + decryptor.finalize() # remove padding unpadder = PKCS7(128).unpadder() clear_text = unpadder.update(padded_text) + unpadder.finalize() return clear_text.decode() def decryptDescription(rec): rec["updated_description"] = decrypt_pbe_with_hmac_sha512_aes_256(rec["description"]) del rec["description"] return rec args = getResolvedOptions(sys.argv, ["JOB_NAME"]) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args["JOB_NAME"], args) # Script generated for node node1 = glueContext.create_dynamic_frame.from_catalog(...) mapped_dyF = Map.apply(frame = node1, f = decryptDescription) # Script generated for node ApplyMapping ApplyMapping_node2 = ApplyMapping.apply(...) # Script generated for node node3 = glueContext.write_dynamic_frame.from_catalog(...) job.commit()
问题根因
该问题不属于AWS Glue Python版本的固有局限性,也不需要强制切换到Scala + 预打包Jar的实现方式。
报错核心原因是:cryptography库的部分内部模块属于_ModuleWithDeprecations类型,无法被Spark的Py4J序列化框架pickle。代码在顶层导入cryptography相关子模块后,Map.apply执行时需要将自定义函数序列化分发到各个Executor节点,序列化过程会扫描函数引用的顶层作用域对象,把不可序列化的cryptography模块实例一并纳入序列化范围,最终触发报错。本地运行为单进程模式,不存在分布式序列化分发步骤,因此可以正常运行。
修复方案
- 调整依赖导入位置:将所有
cryptography相关的导入语句从脚本顶层移动到decrypt_pbe_with_hmac_sha512_aes_256函数内部。修改后的函数结构参考如下:
此时自定义函数的顶层引用不会持有cryptography模块对象,序列化阶段不会尝试序列化这些不可序列化实例,Executor执行函数时会在本地进程内完成依赖导入,从根源规避pickle错误。def decrypt_pbe_with_hmac_sha512_aes_256(obj: str) -> str: # 函数内部导入cryptography依赖,避免顶层不可序列化对象被扫描 from base64 import b64decode from cryptography.hazmat.primitives import hashes from cryptography.hazmat.primitives.kdf.pbkdf2 import PBKDF2HMAC from cryptography.hazmat.backends import default_backend from cryptography.hazmat.primitives.padding import PKCS7 from cryptography.hazmat.primitives.ciphers import Cipher, algorithms, modes # 原有解密逻辑保持不变 - 若需要全局使用cryptography相关能力,可将cryptography依赖打包为独立的zip包,通过Glue作业的
--extra-py-files参数加载,同样需要避免在脚本顶层作用域持有cryptography的模块实例引用。 - 通用规避规则:Glue中所有分布式执行的UDF内用到的第三方库,尽量在函数内部完成导入,可规避绝大多数分布式场景下的序列化报错。
内容的提问来源于stack exchange,提问作者Ryan Shea
相关产品推荐
相关产品推荐

