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

求助: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()

问题原因分析

  1. 分布式执行权限缺失:foreach在Spark Worker节点执行,若Glue Job的IAM角色没有S3的s3:PutObject权限,Worker无法写入S3;部分场景下Worker节点的boto3也无法自动获取AWS凭证。
  2. 日志不可见:Worker节点的print输出不会同步到Glue Job主日志,即使执行出错也无法排查问题。
  3. 低效客户端初始化:每行都初始化boto3客户端,不仅性能差,还可能触发连接限制。
  4. 未利用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 07:32:33