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

如何通过Boto3将S3存储桶文件导入PySpark DataFrame(EMR环境)

Great question! Since you're working in an EMR environment, you actually have a straightforward, effective way to pull multiple file types from S3 into a single PySpark DataFrame—combining Boto3 for file discovery with PySpark's built-in readers works perfectly for this use case. Let me walk you through the process step by step:

Approach: Boto3 + PySpark for Multi-File Type Batch Import

Step 1: Use Boto3 to List and Categorize S3 Files

First, we'll use Boto3 to scan your target S3 bucket (or specific prefix) and group files by their extension (CSV, JSON, etc.). Note: Make sure your EMR instance's IAM profile has s3:ListBucket and s3:GetObject permissions for the bucket you're targeting.

import boto3
from pyspark.sql import SparkSession

# Initialize Spark Session (skip if already running in your EMR environment)
spark = SparkSession.builder.appName("S3MultiTypeImport").getOrCreate()

# Initialize Boto3 S3 client (EMR handles credentials via instance profile automatically)
s3_client = boto3.client('s3')

# Define your target bucket and optional prefix to narrow down files
bucket_name = "your-target-bucket"
prefix = "data/raw-files/"  # Leave empty to scan the entire bucket

# List all non-directory objects in the bucket/prefix
response = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=prefix)
s3_file_keys = [obj['Key'] for obj in response.get('Contents', []) if not obj['Key'].endswith('/')]

# Categorize files by their extension
csv_s3_paths = [f"s3://{bucket_name}/{key}" for key in s3_file_keys if key.lower().endswith('.csv')]
json_s3_paths = [f"s3://{bucket_name}/{key}" for key in s3_file_keys if key.lower().endswith('.json')]
# Add more lists for other types (e.g., parquet_files) if needed

Step 2: Read Each File Type into PySpark DataFrames

Next, use PySpark's optimized readers for each file type. Adjust the options to match your file's structure (e.g., header presence for CSV, multi-line JSON format).

# Read CSV files (customize options like header/sep based on your data)
if csv_s3_paths:
    df_csv = spark.read.csv(
        csv_s3_paths,
        header=True,  # Set to False if your CSV has no header row
        inferSchema=True,  # Replace with a defined StructSchema for better performance with large datasets
        sep=","
    )
else:
    # Create empty DataFrame if no CSV files exist
    df_csv = spark.createDataFrame([], schema=None)

# Read JSON files (adjust multiLine flag based on your JSON structure)
if json_s3_paths:
    df_json = spark.read.json(
        json_s3_paths,
        multiLine=False,  # Set to True if your JSON files are single objects (not line-delimited)
        inferSchema=True
    )
else:
    df_json = spark.createDataFrame([], schema=None)

Step 3: Align Schemas and Union into a Single DataFrame

The critical step here is ensuring all DataFrames have matching schemas before combining them. If your CSV and JSON files have different column sets, you'll need to standardize columns first.

# Get all unique columns across both DataFrames
all_columns = list(set(df_csv.columns + df_json.columns))

# Add missing columns to each DataFrame with null values to align schemas
df_csv_aligned = df_csv.select(*all_columns)
df_json_aligned = df_json.select(*all_columns)

# Union the aligned DataFrames (use unionByName to avoid column order issues)
final_combined_df = df_csv_aligned.unionByName(df_json_aligned, allowMissingColumns=True)

# Verify the result
final_combined_df.show(5)
final_combined_df.printSchema()

Bonus: Skip Boto3 (If You Don't Need Fine-Grained Control)

If you don't need to filter files beyond their extension, PySpark can directly read multiple file types using wildcards—this skips the Boto3 step entirely:

# Read all CSV and JSON files in a prefix
df_csv = spark.read.csv("s3://your-target-bucket/data/raw-files/*.csv", header=True, inferSchema=True)
df_json = spark.read.json("s3://your-target-bucket/data/raw-files/*.json", inferSchema=True)

# Align schemas and union as shown in Step 3

Key Tips for Success

  • Schema Explicitness: For large datasets, define a schema with StructType instead of using inferSchema=True—this drastically improves read performance.
  • Column Alignment: If columns have different names across file types, use withColumnRenamed to standardize them before unioning.
  • Permissions: Double-check your EMR instance profile has the required S3 permissions; if you're using temporary credentials, pass them to Boto3 explicitly via boto3.client('s3', aws_access_key_id='YOUR_KEY', aws_secret_access_key='YOUR_SECRET').
  • Other File Types: Extend this pattern to Parquet (spark.read.parquet()), XML (using the spark-xml library), or any other format supported by PySpark.

内容的提问来源于stack exchange,提问作者Paul Bendevis

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:39:35