如何在AWS Glue PySpark脚本中提取Job ID并新增对应列?
AWS Glue添加Glue Job ID列的解决方案
我是AWS Glue新手,正在处理S3中已被爬虫(Crawler)编目的CSV文件,需要重命名列名、添加新列后将输出以JSON格式写入S3桶。目前已成功为所有记录添加包含当前日期的AcusitionDateTime列,但不知道如何以同样方式添加Glue Job ID列。
现有代码
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 sc = SparkContext.getOrCreate() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) from datetime import datetime import sys from awsglue.utils import getResolvedOptions # Adds a AcusitionDateTime column containing today's date to each record of # the input data set. def AddDateCol(r): r["AcusitionDateTime"] = datetime.now() return r # Script generated for node Data Catalog table datasource0 = glueContext.create_dynamic_frame.from_catalog( database = "test_database2", table_name = "sourcedata_csv", transformation_ctx = "datasource0") # Apply the function to each record of the Dynamic DataFrame datasource0 = Map.apply(frame = datasource0, f = AddDateCol) # Script generated for node ApplyMapping ApplyMapping1 = ApplyMapping.apply( frame=datasource0, mappings=[ ("accountid", "long", "Accountid", "long"), ("accounttype", "string", "Accounttype", "string"), ("accountname", "string", "Accountname", "string"), ("nickname", "string", "Nickname", "string"), ("accounttypeid", "long", "Accounttypeid", "long"), ("acusitiondatetime", "timestamp", "AcusitionDateTime", "timestamp"), ], transformation_ctx="ApplyMapping1" ) # Script generated for node S3 bucket S3bucket_node3 = glueContext.write_dynamic_frame.from_options( frame=ApplyMapping1, connection_type="s3", format="json", connection_options={ "path": "s3://data-lake/dataset/", "compression": "gzip", "partitionKeys": [], }, transformation_ctx="S3bucket_node3", ) job.commit()
期望输出JSON
{ "Accountid": "1234", "Accounttype": "30", "Accountname": "joins", "Nickname": "leejones", "Accounttypeid": "324566", "AcusitionDateTime": "12-01-2023", "Glue_Job_Id": "225273h37dh7dh3w7" }
解决方案
当然可以提取Glue Job ID并添加为新列,具体修改如下:
步骤1:获取Glue作业运行ID
Glue作业运行时会自动传入JOB_RUN_ID参数,通过getResolvedOptions即可获取该值,这就是你需要的Glue_Job_Id。
步骤2:修改列添加逻辑
将日期列和Job ID列的添加逻辑合并到一个函数中,减少转换操作次数。
步骤3:更新列映射
在ApplyMapping中新增Glue_Job_Id的映射规则,确保列名和类型正确。
修改后的完整代码
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 datetime import datetime sc = SparkContext.getOrCreate() glueContext = GlueContext(sc) spark = glueContext.spark_session # 获取作业参数,提取JOB_RUN_ID args = getResolvedOptions(sys.argv, ['JOB_NAME', 'JOB_RUN_ID']) job = Job(glueContext) job.init(args['JOB_NAME'], args) glue_job_id = args['JOB_RUN_ID'] # 同时添加日期列和Glue Job ID列 def AddDateAndJobIdCol(r): r["AcusitionDateTime"] = datetime.now() r["Glue_Job_Id"] = glue_job_id return r # 读取编目后的数据源 datasource0 = glueContext.create_dynamic_frame.from_catalog( database = "test_database2", table_name = "sourcedata_csv", transformation_ctx = "datasource0") # 应用列添加函数 datasource0 = Map.apply(frame = datasource0, f = AddDateAndJobIdCol) # 应用列映射,新增Glue_Job_Id的映射 ApplyMapping1 = ApplyMapping.apply( frame=datasource0, mappings=[ ("accountid", "long", "Accountid", "long"), ("accounttype", "string", "Accounttype", "string"), ("accountname", "string", "Accountname", "string"), ("nickname", "string", "Nickname", "string"), ("accounttypeid", "long", "Accounttypeid", "long"), ("acusitiondatetime", "timestamp", "AcusitionDateTime", "timestamp"), ("glue_job_id", "string", "Glue_Job_Id", "string") ], transformation_ctx="ApplyMapping1" ) # 写入JSON到S3 S3bucket_node3 = glueContext.write_dynamic_frame.from_options( frame=ApplyMapping1, connection_type="s3", format="json", connection_options={ "path": "s3://data-lake/dataset/", "compression": "gzip", "partitionKeys": [], }, transformation_ctx="S3bucket_node3", ) job.commit()
关键修改说明
- 移除了重复的
import语句,优化代码结构 - 通过
getResolvedOptions获取JOB_RUN_ID,这是当前Glue作业运行的唯一标识 - 合并日期列和Job ID列的添加逻辑,提升处理效率
- 在
ApplyMapping中新增Glue_Job_Id的映射,确保输出列名和类型符合预期
内容的提问来源于stack exchange,提问作者Abiodun Adeoye
相关产品推荐
相关产品推荐

