AWS Glue转换XML到JSON遇随机文件名及0KB文件等问题咨询
AWS Glue XML转JSON问题解答
一、随机文件名与0KB文件问题
问题原因
Spark作为分布式计算框架,写入文件时会按分区生成多个part-*开头的文件,这是默认行为。0KB文件通常源于某个分区无数据(比如源数据分区后产生空分区),或是作业执行时生成了空输出分区。
解决方法
- 合并分区生成单个文件(适合小数据量场景):
在写入前用coalesce(1)将DataFrame合并为一个分区,避免多文件输出。注意大数据量下不要使用,会导致单节点负载过高:
若需要自定义文件名,可先写入临时目录,再通过boto3重命名最终文件:# 修改写入代码 source_df.coalesce(1).write.mode("overwrite").json(f"s3://{target_bucket}/{output_file}")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/") - 过滤空数据避免空分区:
若源数据存在大量空行,可先过滤后再写入: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会自动记录上次作业处理的文件/数据范围,下次运行仅处理新增或修改的文件。
配置与代码调整:
- 在Glue作业配置页,将Job bookmark设置为
Enable; - 读取数据时指定
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维护已处理文件的元数据:
- 每次运行作业时,列出S3源桶的所有文件,获取
LastModified时间; - 对比DynamoDB中的记录,筛选出新增/修改的文件;
- 仅处理这些文件,并更新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
相关产品推荐
相关产品推荐

