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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 09:06:39