如何在AWS Glue worker的map函数内实现日志输出?
AWS Glue Worker节点正确打印日志的方法
错误原因
你遇到的PicklingError报错本质是driver端初始化的Glue Logger实例包含不可序列化的_thread.RLock线程锁对象,Spark需要把映射函数序列化后分发到各worker节点执行,不可序列化的logger对象随函数一起传输时就会触发该错误。
正确实现方案
核心逻辑是不要在driver端创建logger后传递给worker执行的函数,改为在worker侧执行的函数内部初始化logger,无需跨节点序列化即可正常输出日志。
示例代码修改
sc = SparkContext() glueContext = GlueContext(sc) # driver端的logger仅在driver执行的代码中使用 driver_logger = glueContext.get_logger() driver_logger.info("starting glue job...") # 正常输出 ... def transform(item): # worker侧执行的函数内部初始化logger from awsglue.context import GlueContext from pyspark.context import SparkContext # 仅在worker节点本地初始化,不需要序列化传输 worker_glue_context = GlueContext(SparkContext.getOrCreate()) worker_logger = worker_glue_context.get_logger() worker_logger.info("starting transform...") # 正常输出 # 此处写转换逻辑 return item Map.apply(frame = dynamicFrame, f = transform)
补充说明
- 你之前配置的Glue连续日志规则对worker侧日志同样生效,worker输出的日志会自动归集到CloudWatch对应的日志组下,和driver日志分开存储,你可以通过日志流前缀区分driver和不同worker节点的日志。
- 建议优先使用Glue自带的logger而非print语句输出日志,logger会自动附带上作业ID、执行节点ID等元信息,后续排查问题时检索效率更高。
- 如果需要在多个worker侧函数中复用logger,可以封装成独立的工具函数,第一次调用时完成初始化,后续调用直接复用已有实例即可,避免重复初始化的性能损耗。
内容的提问来源于stack exchange,提问作者Xiqiang Lin
相关产品推荐
相关产品推荐

