求助:AWS Glue Job无法将PySpark DataFrame每行转为JSON文档
AWS Glue Job实现每行DataFrame存为独立S3 JSON文件的问题解决
我尝试用AWS Glue Job将PySpark DataFrame的每一行保存为独立的JSON文件到S3,但功能未正常生效。试过在函数内初始化boto3客户端、转RDD处理等方法都没用,求解决办法。附上我的代码:
import sys import boto3 import json from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from pyspark.sql import Row from datetime import datetime # Initialize GlueContext, Spark session, and Job sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) args = { 'JOB_NAME': 'SampleGlueJob' } job.init(args['JOB_NAME'], args) # Define the S3 bucket and folder to save the JSON files s3_bucket = "<bucket name>" s3_folder = "<Folder name>" # Create a dummy PySpark DataFrame data = [ Row(id=1, name="Alice", age=29), Row(id=2, name="Bob", age=34), Row(id=3, name="Charlie", age=25) ] df = spark.createDataFrame(data) # Show the DataFrame print("Dummy DataFrame:") df.show() # Function to save each row as a JSON file in S3 def save_row_as_json(row): try: # Initialize the S3 client inside the function s3 = boto3.client('s3') row_dict = row.asDict() unique_id = row_dict.get("id") # Use a unique identifier for the file name current_date = datetime.now().strftime("%Y%m%d%H%M%S") filename = f"{unique_id}_{current_date}.json" json_body = json.dumps(row_dict, indent=2) json_file_path = f"{s3_folder}/{filename}" # Save to S3 s3.put_object(Bucket=s3_bucket, Key=json_file_path, Body=json_body) print(f"Successfully saved row {unique_id} to s3://{s3_bucket}/{json_file_path}") except Exception as e: print(f"Error saving row {row_dict.get('id')} to S3: {str(e)}") # Save each row as a separate JSON file df.foreach(save_row_as_json) # Commit the Glue job job.Commit()
问题原因分析
- 分布式执行权限缺失:
foreach在Spark Worker节点执行,若Glue Job的IAM角色没有S3的s3:PutObject权限,Worker无法写入S3;部分场景下Worker节点的boto3也无法自动获取AWS凭证。 - 日志不可见:Worker节点的
print输出不会同步到Glue Job主日志,即使执行出错也无法排查问题。 - 低效客户端初始化:每行都初始化boto3客户端,不仅性能差,还可能触发连接限制。
- 未利用Spark原生能力:手动用boto3写入不是Spark最优方案,原生API更稳定高效。
解决方案
方案一:修复现有代码(支持自定义文件名)
关键改进
- 改用
foreachPartition,每个分区初始化一次boto3客户端,减少开销 - 使用Spark日志系统记录错误,确保能在CloudWatch中查看
- 显式指定AWS区域,避免节点区域不匹配
- 确保IAM角色拥有S3写入权限
修复后代码
import sys import boto3 import json from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from pyspark.sql import Row from datetime import datetime import logging # 获取Spark日志记录器 logger = logging.getLogger(__name__) # Initialize GlueContext, Spark session, and Job sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) args = { 'JOB_NAME': 'SampleGlueJob' } job.init(args['JOB_NAME'], args) # Define the S3 bucket and folder to save the JSON files s3_bucket = "<bucket name>" s3_folder = "<Folder name>" aws_region = "us-east-1" # 替换为你的AWS区域 # Create a dummy PySpark DataFrame data = [ Row(id=1, name="Alice", age=29), Row(id=2, name="Bob", age=34), Row(id=3, name="Charlie", age=25) ] df = spark.createDataFrame(data) # Show the DataFrame print("Dummy DataFrame:") df.show() # 每个分区初始化一次S3客户端,批量处理分区内的行 def save_partition_as_json(partition): try: # 初始化S3客户端,指定区域 s3 = boto3.client('s3', region_name=aws_region) for row in partition: row_dict = row.asDict() unique_id = row_dict.get("id") current_date = datetime.now().strftime("%Y%m%d%H%M%S") filename = f"{unique_id}_{current_date}.json" json_body = json.dumps(row_dict, indent=2) json_file_path = f"{s3_folder}/{filename}" s3.put_object(Bucket=s3_bucket, Key=json_file_path, Body=json_body) logger.info(f"Successfully saved row {unique_id} to s3://{s3_bucket}/{json_file_path}") except Exception as e: logger.error(f"Error saving partition rows: {str(e)}", exc_info=True) # 用foreachPartition替代foreach df.rdd.foreachPartition(save_partition_as_json) # Commit the Glue job job.commit()
方案二:Spark原生写入(高效稳定,无需自定义文件名)
如果不需要自定义文件名,直接用Spark原生API实现每行一个文件:
# 省略初始化代码,直接使用df写入 df.write.mode("overwrite") \ .option("maxRecordsPerFile", 1) \ .json(f"s3://{s3_bucket}/{s3_folder}")
- 优势:Spark自动处理分布式写入、重试逻辑,性能更优
- 劣势:文件名由Spark自动生成(如
part-00000-xxxx.json),无法自定义
方案三:自定义文件名的高效写入(结合RDD与Hadoop API)
追求性能且需要自定义文件名时,使用Hadoop FileSystem API:
from pyspark.sql.functions import col, concat_ws, lit, to_json from org.apache.hadoop.fs import Path from org.apache.hadoop.io import BytesWritable, Text from org.apache.hadoop.mapred import SequenceFileOutputFormat # 生成自定义文件名和JSON内容 df_with_filename = df.withColumn( "filename", concat_ws("_", col("id"), lit(datetime.now().strftime("%Y%m%d%H%M%S")), lit(".json")) ).withColumn( "json_content", to_json(col("*")) ) # 转成RDD:(文件名, JSON内容) rdd = df_with_filename.rdd.map(lambda x: (Text(x.filename), BytesWritable(x.json_content.encode('utf-8')))) # 保存到S3 rdd.saveAsHadoopFile( f"s3://{s3_bucket}/{s3_folder}", SequenceFileOutputFormat, Text, BytesWritable )
必查事项
- 确认Glue Job的IAM角色包含以下权限(替换为你的S3路径):
{ "Effect": "Allow", "Action": "s3:PutObject", "Resource": "arn:aws:s3:::<bucket name>/<Folder name>/*" } - 查看CloudWatch Logs中的Glue Job日志,排查是否有权限错误或其他异常
内容的提问来源于stack exchange,提问作者kiran u
相关产品推荐
相关产品推荐

