如何使用S3分段上传实现CSV数据流式上传以降低内存占用?
Great question—you’re exactly right about the memory bottleneck here. Loading an entire CSV (especially a large one pulled from a database) into memory before uploading to S3 can quickly eat up your resources, and S3’s multipart upload feature is made for this exact scenario. Instead of buffering everything locally, you can stream chunks of your CSV directly to S3 as you generate them, keeping memory usage low.
Here’s a step-by-step approach with code examples using boto3:
Core Idea
We’ll:
- Initialize an S3 multipart upload session
- Pull chunks of data from your database (instead of all at once)
- Convert each chunk to CSV format
- Upload each CSV chunk as a separate part to S3
- Finalize the multipart upload once all parts are sent
Code Implementation
First, set up your S3 client and initialize the multipart upload:
import boto3 import csv from io import StringIO # Initialize S3 client s3 = boto3.client('s3') BUCKET_NAME = "your-target-bucket" OBJECT_KEY = "path/to/your/file.csv" # Start a multipart upload session mpu_response = s3.create_multipart_upload(Bucket=BUCKET_NAME, Key=OBJECT_KEY) upload_id = mpu_response['UploadId'] parts = [] part_number = 1
Next, pull data in chunks (replace the mock DB chunk function with your actual database read logic) and upload each chunk:
# Mock function: Replace this with your database pagination/streaming logic def get_db_data_chunks(chunk_size=10_000): """Yield chunks of data from your database instead of loading all at once.""" # Example: Simulate fetching 10 chunks of 10k rows each for chunk_idx in range(10): start_row = chunk_idx * chunk_size end_row = start_row + chunk_size yield [{"id": row, "data": f"sample_data_{row}"} for row in range(start_row, end_row)] # Process each data chunk and upload as an S3 part for data_chunk in get_db_data_chunks(): # Write the chunk to a CSV buffer csv_buffer = StringIO() writer = csv.DictWriter(csv_buffer, fieldnames=["id", "data"]) # Write CSV header only on the first chunk if part_number == 1: writer.writeheader() writer.writerows(data_chunk) csv_buffer.seek(0) # Reset buffer position to start # Upload the chunk as an S3 part upload_response = s3.upload_part( Bucket=BUCKET_NAME, Key=OBJECT_KEY, UploadId=upload_id, PartNumber=part_number, Body=csv_buffer.getvalue().encode("utf-8") ) # Track part details for finalization parts.append({"PartNumber": part_number, "ETag": upload_response["ETag"]}) part_number += 1 csv_buffer.close()
Finally, complete the multipart upload (and add error handling to abort if something goes wrong):
try: # Finalize the multipart upload to combine all parts into a single file s3.complete_multipart_upload( Bucket=BUCKET_NAME, Key=OBJECT_KEY, UploadId=upload_id, MultipartUpload={"Parts": parts} ) print(f"Successfully uploaded {OBJECT_KEY} to S3!") except Exception as e: # If anything fails, abort the multipart upload to avoid leftover parts (which cost money) s3.abort_multipart_upload( Bucket=BUCKET_NAME, Key=OBJECT_KEY, UploadId=upload_id ) print(f"Upload failed, aborted multipart session: {str(e)}") raise
Key Notes & Best Practices
- Part Size Requirements: S3 requires each part (except the last one) to be at least 5MB. Adjust your
chunk_sizeso the resulting CSV chunk meets this threshold—you can calculate based on your average row size, or use a byte-based buffer instead of row-count-based chunks. - Automatic Multipart with
upload_fileobj: If you prefer a more hands-off approach, boto3’supload_fileobjcan automatically handle multipart uploads if you set amultipart_thresholdin the config. This works great if you can stream your CSV data directly as a file-like object:from botocore.config import Config # Configure auto-multipart for files over 10MB s3_config = Config(multipart_threshold=10 * 1024 * 1024) s3 = boto3.client('s3', config=s3_config) # Stream CSV data directly to S3 def csv_stream_generator(): # Yield CSV header first yield b"id,data\n" # Yield each data chunk as CSV bytes for data_chunk in get_db_data_chunks(): buffer = StringIO() writer = csv.DictWriter(buffer, fieldnames=["id", "data"]) writer.writerows(data_chunk) yield buffer.getvalue().encode("utf-8") buffer.close() s3.upload_fileobj(csv_stream_generator(), BUCKET_NAME, OBJECT_KEY) - Error Handling: Always abort incomplete multipart uploads—S3 charges for stored parts until you either complete or abort the session.
内容的提问来源于stack exchange,提问作者James

