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

寻求S3上连接两个Parquet表并输出至S3的低成本解决方案

Cost-Effective Solution to Join Parquet Tables on S3 & Write Back

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

  1. 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()
  1. 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/")
  1. 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.partitions if 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

  1. 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/';
  1. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:04:34