AWS Glue Job输出文件批量重命名实现方案问询
Got it, let's tackle this problem step by step. Since AWS Glue doesn't let you name output files directly during job execution, we can add a post-processing step in your Glue Job to batch rename the output files to match their corresponding input filenames. We'll use Python (PySpark) since it's fully compatible with Glue Jobs, paired with boto3 for S3 operations.
Here's a complete code example you can integrate into your existing Glue Job (add this at the end of your ETL logic, after the output has been written to S3):
import boto3 from pyspark.context import SparkContext from awsglue.context import GlueContext import pyspark.sql.functions as F # Initialize contexts (skip if already set up in your job) sc = SparkContext.getOrCreate() glueContext = GlueContext(sc) # Configure your S3 paths - update these with your actual bucket/prefixes INPUT_BUCKET = "your-input-bucket" INPUT_PREFIX = "raw-data/" # Leave empty if files are at bucket root OUTPUT_BUCKET = "your-output-bucket" OUTPUT_PREFIX = "processed-data/" # Path where Glue writes part-* files # Initialize S3 client s3 = boto3.client('s3') def get_input_filenames(): """Fetch all input filenames (excluding folder entries)""" response = s3.list_objects_v2(Bucket=INPUT_BUCKET, Prefix=INPUT_PREFIX) input_files = [] for obj in response.get('Contents', []): # Skip folder paths (keys ending with /) if not obj['Key'].endswith('/'): # Extract just the filename (e.g., "user_data.csv" from "raw-data/user_data.csv") filename = obj['Key'].split('/')[-1] input_files.append(filename) return input_files def get_output_part_files(): """Fetch all Glue-generated part files in the output directory""" response = s3.list_objects_v2(Bucket=OUTPUT_BUCKET, Prefix=OUTPUT_PREFIX) part_files = [] for obj in response.get('Contents', []): # Target Glue's default part file naming pattern if obj['Key'].startswith(f"{OUTPUT_PREFIX}part-") and not obj['Key'].endswith('/'): part_files.append(obj['Key']) return part_files def rename_output_files(input_filenames, part_files): """Rename part files to match input filenames. Adjust logic for 1:many mappings if needed""" # Assuming 1 input file maps to 1 output part file (modify if you have multiple parts per input) for idx, target_filename in enumerate(input_filenames): if idx >= len(part_files): print(f"No matching part file for input: {target_filename}") continue source_key = part_files[idx] target_key = f"{OUTPUT_PREFIX}{target_filename}" # Copy to new name (S3 doesn't have true rename, so copy + delete) s3.copy_object( Bucket=OUTPUT_BUCKET, CopySource={'Bucket': OUTPUT_BUCKET, 'Key': source_key}, Key=target_key ) # Delete original part file s3.delete_object(Bucket=OUTPUT_BUCKET, Key=source_key) print(f"Successfully renamed: {source_key} → {target_key}") # Run the rename workflow input_files = get_input_filenames() output_parts = get_output_part_files() rename_output_files(input_files, output_parts)
Key Notes for the Code:
- IAM Permissions: Ensure your Glue Job's IAM role has these permissions:
s3:ListBucketon both input and output bucketss3:GetObject,s3:PutObject,s3:DeleteObjecton the output bucket
- Input-Output Mapping: The code uses a simple index-based 1:1 mapping. If your job produces multiple part files per input, modify the logic (e.g., use file creation timestamps or track input filenames during ETL with
F.input_file_name()). - Overwrite Protection: Add a check with
s3.head_object()before copying if you want to avoid overwriting existing files with the target name.
To run this job automatically on a schedule:
- Go to the AWS Glue Console > Jobs > Select your job
- Click Add trigger > Choose Schedule
- Configure your frequency (e.g., daily, hourly) using a cron expression or predefined interval
- Save the trigger — your job will now run at the specified times
If index-based mapping isn't reliable, capture input filenames while reading data to create a clear mapping:
# When reading input data, add a column with the source filename df = glueContext.create_dynamic_frame.from_options( connection_type="s3", connection_options={"paths": [f"s3://{INPUT_BUCKET}/{INPUT_PREFIX}"]}, format="csv" ).toDF() # Add input filename column df = df.withColumn("source_filename", F.split(F.input_file_name(), "/").getItem(-1))
You can use this column to partition output data by filename (creating subfolders) or store the mapping in a metadata table to reference during renaming.
内容的提问来源于stack exchange,提问作者PFGX

