如何在AWS Step Functions中传递Glue Job输出至下一个Glue Job
问题分析
你当前的Step Functions配置无法获取Glue Job1的输出,核心原因是:arn:aws:states:::glue:startJobRun.sync返回的是Glue Job运行的元数据(如JobRunId、运行状态),而非你脚本中print的内容。因此$.job1Result.output_value这个路径不存在,导致Job2无法拿到正确输入。
解决方案
推荐采用Glue Job将结果写入S3,Step Functions读取后传递给下一个Job的方案,这是最稳定可靠的实现方式。
1. 修改Glue Job1脚本:将结果写入S3
把原来仅打印结果的逻辑,改为将结果序列化后写入S3存储桶(替换为你自己的桶和路径):
import sys import json import boto3 from awsglue.context import GlueContext from pyspark.context import SparkContext # 初始化上下文 sparkContext = SparkContext() glueContext = GlueContext(sparkContext) s3_client = boto3.client('s3') # 获取Step Functions传入的参数 input_value = sys.argv[1] # 处理逻辑 output_value = input_value + "pass" result = {"output_value": output_value} result_json = json.dumps(result) # 将结果写入S3 bucket = "your-custom-output-bucket" file_path = "glue-job-outputs/job1-result.json" s3_client.put_object(Bucket=bucket, Key=file_path, Body=result_json) # 可选:保留打印方便调试 print(result_json)
2. 更新Step Functions流程:新增读取S3的步骤
在Job1执行完成后,添加一个任务读取S3中的Job1结果,再传递给Job2:
{ "StartAt": "Job1", "States": { "Job1": { "Type": "Task", "Resource": "arn:aws:states:::glue:startJobRun.sync", "Parameters": { "JobName": "Job1", "Arguments": { "input.$": "$.input" } }, "Next": "FetchJob1Result" }, "FetchJob1Result": { "Type": "Task", "Resource": "arn:aws:states:::s3:getObject", "Parameters": { "Bucket": "your-custom-output-bucket", "Key": "glue-job-outputs/job1-result.json" }, "ResultPath": "$.job1Result", "Next": "Job2" }, "Job2": { "Type": "Task", "Resource": "arn:aws:states:::glue:startJobRun.sync", "Parameters": { "JobName": "Job2", "Arguments": { "input.$": "$.job1Result.Body.output_value" } }, "ResultPath": "$.job2Result", "End": true } } }
3. 权限配置(必做)
确保以下权限配置正确:
- 执行Step Functions的角色需要拥有
s3:GetObject权限,允许读取目标S3桶的内容 - Glue Job的执行角色需要拥有
s3:PutObject权限,允许写入目标S3桶
可选优化:将S3输出路径作为参数传入Glue Job
如果需要动态指定输出路径,可以在Step Functions中把S3路径作为参数传给Job1,避免脚本中硬编码:
// 在Job1的Parameters中添加 "Arguments": { "input.$": "$.input", "--s3_output_path": "s3://your-custom-output-bucket/glue-job-outputs/job1-result.json" }
然后在Job1脚本中读取这个参数:
s3_output_path = sys.argv[2] # 解析bucket和key from urllib.parse import urlparse parsed = urlparse(s3_output_path) bucket = parsed.netloc file_path = parsed.path.lstrip('/')
内容的提问来源于stack exchange,提问作者Propsz
相关产品推荐
相关产品推荐

