寻求S3上连接两个Parquet表并输出至S3的低成本解决方案
Hey there! Let's walk through the most budget-friendly ways to join your two Parquet tables stored on S3, considering their massive size difference (10.5 TiB vs. 200 GiB). The core trick here is to avoid shuffling the huge table_a—shuffling large datasets eats up compute time and costs, so we’ll leverage broadcasting the smaller table_b instead.
Core Strategy: Broadcast the Small Table
Since table_b is only 200 GiB (way smaller than table_a), we can broadcast it to all worker nodes. This lets each node join its chunk of table_a with a local copy of table_b, eliminating the need to move table_a’s data across the network—this is the biggest cost saver here.
Option 1: Apache Spark (Flexible & Cost-Efficient with Spot Instances)
Spark has top-tier support for Parquet and S3, and it’s easy to configure for broadcast joins. For long-running or recurring jobs, using EC2 Spot Instances can cut compute costs by up to 90% compared to On-Demand instances.
Step-by-Step Implementation
- Initialize Spark Session with broadcast threshold adjustments (to accommodate table_b’s size):
from pyspark.sql import SparkSession from pyspark.sql.functions import broadcast # Configure Spark to allow broadcasting tables up to 200 GiB (214748364800 bytes) spark = SparkSession.builder \ .appName("S3ParquetJoin") \ .config("spark.sql.autoBroadcastJoinThreshold", "214748364800") \ .getOrCreate()
- Read Tables from S3:
table_a = spark.read.parquet("s3://your-bucket/path/to/table_a/") table_b = spark.read.parquet("s3://your-bucket/path/to/table_b/")
- Perform Broadcast Join & Write Result:
# Join on the matching key (adjust the join condition to fit your actual business logic) joined_df = table_a.join(broadcast(table_b), on="name", how="inner") # Write back to S3 with Snappy compression (reduces storage costs and future read times) joined_df.write \ .mode("overwrite") \ .option("compression", "snappy") \ .parquet("s3://your-bucket/path/to/table_c/")
Spark Cost Optimization Tips
- Use EC2 Spot Instances for your Spark cluster—AWS offers deep discounts for unused cloud capacity.
- Tune parallelism: Ensure each task processes 128–256 MB of data (adjust
spark.sql.shuffle.partitionsif needed) to avoid underutilizing resources. - Enable S3 lifecycle rules to archive old data or delete temporary files after the job finishes.
Option 2: AWS Athena (Serverless, No Cluster Management)
Athena is perfect for one-off or ad-hoc joins—you pay only for the data scanned, and there’s no infrastructure to manage. It automatically optimizes joins when possible, but we can explicitly hint to broadcast table_b to ensure optimal performance.
Step-by-Step Implementation
- Create External Tables (if you haven’t already) pointing to your S3 Parquet data:
CREATE EXTERNAL TABLE table_a ( id INT, name STRING ) STORED AS PARQUET LOCATION 's3://your-bucket/path/to/table_a/'; CREATE EXTERNAL TABLE table_b ( name STRING, id INT ) STORED AS PARQUET LOCATION 's3://your-bucket/path/to/table_b/';
- Run CTAS Query with Broadcast Hint:
CREATE TABLE table_c WITH ( format = 'PARQUET', parquet_compression = 'SNAPPY', external_location = 's3://your-bucket/path/to/table_c/' ) AS SELECT a.id AS a_id, a.name, b.id AS b_id FROM table_a a /*+ BROADCAST(b) */ -- Force broadcast of the smaller table to avoid shuffling table_a JOIN table_b b ON a.name = b.name;
Athena Cost Optimization Tips
- Use partitioned tables if table_a is partitioned (e.g., by date) to reduce the amount of data scanned per query.
- Enable Athena Result Reuse to avoid re-scanning data if you run similar queries multiple times.
- Set up S3 lifecycle rules for the output location to move data to cheaper storage tiers (like S3 Glacier) after a certain period.
Option 3: AWS Glue (Serverless ETL for Recurring Jobs)
If you need to schedule this join periodically, Glue is a great middle ground—it’s serverless, integrates seamlessly with S3 and Parquet, and uses Spark under the hood. You can configure broadcast joins just like in Spark, and Glue handles cluster provisioning for you.
Quick Example Script
import sys from awsglue.context import GlueContext from pyspark.sql.functions import broadcast from awsglue.dynamicframe import DynamicFrame glueContext = GlueContext(SparkContext.getOrCreate()) # Read tables from S3 table_a = glueContext.create_dynamic_frame.from_options( connection_type="s3", connection_options={"paths": ["s3://your-bucket/path/to/table_a/"]}, format="parquet" ).toDF() table_b = glueContext.create_dynamic_frame.from_options( connection_type="s3", connection_options={"paths": ["s3://your-bucket/path/to/table_b/"]}, format="parquet" ).toDF() # Perform broadcast join joined_df = table_a.join(broadcast(table_b), on="name", how="inner") # Write result back to S3 glueContext.write_dynamic_frame.from_options( frame=DynamicFrame.fromDF(joined_df, glueContext, "joined_df"), connection_type="s3", connection_options={"path": "s3://your-bucket/path/to/table_c/"}, format="parquet", format_options={"compression": "snappy"} )
Which Option Should You Choose?
- One-off job: Go with Athena—no setup, pay-as-you-go, and minimal effort.
- Recurring/complex job: Use Spark with Spot Instances or AWS Glue for better control and lower long-term costs.
- Always prioritize broadcasting the small table to avoid expensive shuffles of table_a’s 10.5 TiB of data.
内容的提问来源于stack exchange,提问作者Thagor

