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

AWS Glue转换XML到JSON遇随机文件名及0KB文件等问题咨询

AWS Glue XML转JSON问题解答

一、随机文件名与0KB文件问题

问题原因

Spark作为分布式计算框架,写入文件时会按分区生成多个part-*开头的文件,这是默认行为。0KB文件通常源于某个分区无数据(比如源数据分区后产生空分区),或是作业执行时生成了空输出分区。

解决方法

  1. 合并分区生成单个文件(适合小数据量场景):
    在写入前用coalesce(1)将DataFrame合并为一个分区,避免多文件输出。注意大数据量下不要使用,会导致单节点负载过高:
    # 修改写入代码
    source_df.coalesce(1).write.mode("overwrite").json(f"s3://{target_bucket}/{output_file}")
    
    若需要自定义文件名,可先写入临时目录,再通过boto3重命名最终文件:
    import boto3
    import os
    
    temp_path = f"s3://{target_bucket}/temp_output"
    final_file = f"{output_file}.json"
    
    # 写入临时目录
    source_df.coalesce(1).write.mode("overwrite").json(temp_path)
    
    # 获取临时目录下的part文件
    s3 = boto3.client('s3')
    response = s3.list_objects_v2(Bucket=target_bucket, Prefix="temp_output/part-")
    part_file = response['Contents'][0]['Key']
    
    # 复制并重命名到最终路径
    s3.copy_object(
        Bucket=target_bucket,
        CopySource=f"{target_bucket}/{part_file}",
        Key=final_file
    )
    
    # 清理临时文件
    s3.delete_object(Bucket=target_bucket, Key=part_file)
    s3.delete_object(Bucket=target_bucket, Key="temp_output/_SUCCESS")
    s3.delete_object(Bucket=target_bucket, Key="temp_output/")
    
  2. 过滤空数据避免空分区:
    若源数据存在大量空行,可先过滤后再写入:
    source_df = source_df.dropna(how='all')
    

二、无需Glue表的转换方式

完全可以跳过Glue Catalog表,直接读取S3上的XML文件,以下是两种常用方案:

方式1:Spark原生XML读取

Glue作业默认包含spark-xml依赖,直接用Spark DataFrame读取,需指定XML的根节点(rowTag):

# 直接读取S3 XML文件,替换your_root_tag为实际XML根节点
source_df = spark.read.format("xml") \
    .option("rowTag", "your_root_tag") \
    .load(f"s3://{source_bucket}/")

# 写入JSON
source_df.write.mode("overwrite").json(f"s3://{target_bucket}/{output_file}")

方式2:Glue DynamicFrame直接读取

使用glueContext.create_dynamic_frame.from_options指定数据源格式与路径,无需Catalog表:

datasource = glueContext.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={"paths": [f"s3://{source_bucket}/"]},
    format="xml",
    format_options={"rowTag": "your_root_tag"}
)

source_df = datasource.toDF()
# 后续写入逻辑不变

这种方式更灵活,适合临时或动态文件处理场景。

三、增量更新(仅处理新增文件)

方法1:启用Glue Job Bookmark

这是官方推荐的增量处理方案,Job Bookmark会自动记录上次作业处理的文件/数据范围,下次运行仅处理新增或修改的文件。

配置与代码调整:

  1. 在Glue作业配置页,将Job bookmark设置为Enable;
  2. 读取数据时指定transformation_ctx,确保Bookmark生效:
    datasource = glueContext.create_dynamic_frame.from_options(
        connection_type="s3",
        connection_options={
            "paths": [f"s3://{source_bucket}/"],
            "recurse": True  # 开启递归读取,Bookmark才能跟踪文件
        },
        format="xml",
        format_options={"rowTag": "your_root_tag"},
        transformation_ctx="datasource"  # 必填参数,关联Bookmark
    )
    

方法2:自定义元数据跟踪

若Job Bookmark不满足需求,可通过DynamoDB维护已处理文件的元数据:

  1. 每次运行作业时,列出S3源桶的所有文件,获取LastModified时间;
  2. 对比DynamoDB中的记录,筛选出新增/修改的文件;
  3. 仅处理这些文件,并更新DynamoDB记录。

示例代码片段:

import boto3

# 初始化DynamoDB客户端
dynamodb = boto3.resource('dynamodb')
table = dynamodb.Table('processed_files')

# 获取S3源桶文件列表
s3 = boto3.client('s3')
response = s3.list_objects_v2(Bucket=source_bucket)
files = response['Contents']

# 筛选未处理的文件
unprocessed_files = []
for file in files:
    file_key = file['Key']
    last_modified = file['LastModified'].isoformat()
    
    # 查询是否已处理
    db_response = table.get_item(Key={'file_key': file_key})
    if 'Item' not in db_response or db_response['Item']['last_modified'] != last_modified:
        unprocessed_files.append(f"s3://{source_bucket}/{file_key}")

# 处理未处理文件并追加写入
if unprocessed_files:
    source_df = spark.read.format("xml") \
        .option("rowTag", "your_root_tag") \
        .load(unprocessed_files)
    
    source_df.write.mode("append").json(f"s3://{target_bucket}/{output_file}")
    
    # 更新DynamoDB记录
    with table.batch_writer() as batch:
        for file in files:
            batch.put_item(
                Item={
                    'file_key': file['Key'],
                    'last_modified': file['LastModified'].isoformat()
                }
            )

内容的提问来源于stack exchange,提问作者Pradeep Charan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 21:54:56