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

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.

1. PySpark Code for Batch Rename in Glue Job

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:ListBucket on both input and output buckets
    • s3:GetObject, s3:PutObject, s3:DeleteObject on 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.
2. Set Up Scheduled Execution

To run this job automatically on a schedule:

  1. Go to the AWS Glue Console > Jobs > Select your job
  2. Click Add trigger > Choose Schedule
  3. Configure your frequency (e.g., daily, hourly) using a cron expression or predefined interval
  4. Save the trigger — your job will now run at the specified times
3. Advanced: Track Input Filenames During ETL

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 08:31:17