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

如何将Glue代码拆分为多个文件并在作业中导入?

在AWS Glue作业中拆分代码并导入自定义模块

一、代码拆分方案

按功能职责拆分文件,让结构更清晰:

  • main.py:Glue作业入口,负责核心流程调度(初始化上下文、解析参数、调用工具函数、执行ETL逻辑)
  • logging.py:封装自定义日志配置,统一日志格式、级别和输出目标
  • utils.py:存放通用工具函数,比如数据校验、S3操作、格式转换等

1. logging.py 示例

import logging
import sys
from awsglue.utils import getResolvedOptions

def get_custom_logger(name):
    # 从作业参数读取日志级别,默认INFO
    args = getResolvedOptions(sys.argv, ['LOG_LEVEL']) if 'LOG_LEVEL' in sys.argv else {}
    log_level = args.get('LOG_LEVEL', 'INFO').upper()
    
    logger = logging.getLogger(name)
    logger.setLevel(getattr(logging, log_level))
    
    # 避免重复添加日志处理器
    if not logger.handlers:
        formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
        console_handler = logging.StreamHandler()
        console_handler.setFormatter(formatter)
        logger.addHandler(console_handler)
    
    return logger

2. utils.py 示例

import boto3

def validate_s3_path(path):
    if not path.startswith('s3://'):
        raise ValueError(f"无效S3路径:{path}")
    return path

def count_s3_files(bucket, prefix):
    s3_client = boto3.client('s3')
    response = s3_client.list_objects_v2(Bucket=bucket, Prefix=prefix)
    return response.get('KeyCount', 0)

3. main.py 示例(作业入口)

import sys
from awsglue.context import GlueContext
from pyspark.context import SparkContext
from awsglue.utils import getResolvedOptions

# 导入自定义模块
from logging import get_custom_logger
from utils import validate_s3_path, count_s3_files

def main():
    # 初始化Spark/Glue上下文
    sc = SparkContext()
    glue_context = GlueContext(sc)
    spark = glue_context.spark_session
    
    # 初始化自定义日志器
    logger = get_custom_logger(__name__)
    logger.info("Glue作业启动,开始执行初始化")
    
    # 解析作业参数
    args = getResolvedOptions(sys.argv, ['INPUT_S3_PATH', 'OUTPUT_S3_PATH'])
    input_path = args['INPUT_S3_PATH']
    output_path = args['OUTPUT_S3_PATH']
    
    # 校验输入输出路径
    try:
        validated_input = validate_s3_path(input_path)
        validated_output = validate_s3_path(output_path)
        logger.info(f"路径校验完成:输入={validated_input},输出={validated_output}")
    except ValueError as e:
        logger.error(f"路径校验失败:{str(e)}")
        sys.exit(1)
    
    # 核心ETL逻辑示例
    df = spark.read.parquet(validated_input)
    logger.info(f"读取数据完成,总行数:{df.count()}")
    
    df.write.mode('overwrite').parquet(validated_output)
    logger.info("数据写入完成,作业执行结束")

if __name__ == '__main__':
    main()

二、Glue作业中导入自定义模块的两种方法

方法1:打包ZIP上传S3,配置Python库路径

  1. 将logging.py、utils.py打包成ZIP包(注意:ZIP根目录直接放.py文件,不要嵌套文件夹)
  2. 把ZIP包上传到S3存储桶,比如s3://your-glue-resources/custom-modules.zip
  3. 在Glue作业配置页,找到「Python库路径」配置项,填入ZIP包的S3路径
  4. 保存配置后,作业运行时Glue会自动将ZIP包解压到Python的sys.path中,main.py可直接导入模块

方法2:同目录存放,直接导入

  1. 将所有.py文件上传到同一个S3目录,比如s3://your-glue-scripts/etl-job/
  2. 在Glue作业配置中,将「脚本路径」设置为main.py的S3路径:s3://your-glue-scripts/etl-job/main.py
  3. main.py中直接导入自定义模块即可(Glue会自动将脚本所在S3目录添加到sys.path)

三、关键注意事项

  • 自定义模块使用的依赖必须是Glue环境预装的(如boto3、logging),若需第三方库,需一起打包到ZIP并配置到Python库路径
  • 日志会自动同步到CloudWatch Logs,可通过作业的「日志组」配置查看
  • 不要在自定义模块中初始化Spark/GlueContext,统一在main.py初始化后传递给工具函数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 05:00:16