Spark SQL 中的分区和分桶是两种重要的数据组织技术,用于优化查询性能。

1. 分区(Partitioning)

1.1 分区概念

分区是根据某个列的值将数据物理分割到不同的目录中,类似于文件系统的文件夹结构。

graph TD
    A[原始表] --> B[按日期分区]
    B --> C[2024-01-01/]
    B --> D[2024-01-02/]
    B --> E[2024-01-03/]
    
    C --> F[data1.parquet]
    C --> G[data2.parquet]
    D --> H[data3.parquet]
    E --> I[data4.parquet]

1.2 分区表创建和使用

-- 创建分区表
CREATE TABLE sales_partitioned (
    product_id STRING,
    amount DECIMAL(10,2),
    customer_id STRING
)
PARTITIONED BY (sale_date DATE, region STRING)
STORED AS PARQUET
LOCATION '/data/sales_partitioned';

-- 插入数据(自动分区)
INSERT INTO sales_partitioned 
PARTITION (sale_date='2024-01-01', region='north')
SELECT product_id, amount, customer_id 
FROM raw_sales 
WHERE sale_date = '2024-01-01' AND region = 'north';

-- 动态分区插入
SET spark.sql.sources.partitionOverwriteMode=dynamic;
INSERT OVERWRITE TABLE sales_partitioned 
PARTITION (sale_date, region)
SELECT product_id, amount, customer_id, sale_date, region
FROM raw_sales;

-- 查询时利用分区剪裁
SELECT * FROM sales_partitioned 
WHERE sale_date BETWEEN '2024-01-01' AND '2024-01-31' 
  AND region = 'north';  -- 只扫描相关分区

1.3 DataFrame API 操作分区

// 创建分区表
val df = spark.read.parquet("/data/raw_sales")

df.write
  .partitionBy("sale_date", "region")
  .format("parquet")
  .mode("overwrite")
  .save("/data/sales_partitioned")

// 读取特定分区
val partitionedDF = spark.read
  .format("parquet")
  .option("basePath", "/data/sales_partitioned")
  .load("/data/sales_partitioned/sale_date=2024-01-01/region=north")

// 查看分区信息
spark.sql("SHOW PARTITIONS sales_partitioned").show()

// 添加新分区
spark.sql("ALTER TABLE sales_partitioned ADD PARTITION (sale_date='2024-01-02', region='south')")

// 删除分区
spark.sql("ALTER TABLE sales_partitioned DROP PARTITION (sale_date='2023-12-31')")

2. 分桶(Bucketing)

2.1 分桶概念

分桶是根据哈希函数将数据分布到固定数量的文件中,每个文件包含哈希值相同的记录。

原始数据
哈希分桶
Bucket 0
Bucket 1
Bucket 2
...
Bucket N-1
user1, userN+1...
user2, userN+2...

2.2 分桶表创建和使用

-- 创建分桶表
CREATE TABLE users_bucketed (
    user_id BIGINT,
    username STRING,
    email STRING,
    registration_date DATE
)
CLUSTERED BY (user_id) INTO 32 BUCKETS
STORED AS PARQUET
LOCATION '/data/users_bucketed';

-- 插入数据(需要启用分桶写入)
SET spark.sql.sources.bucketing.enabled=true;
SET spark.sql.adaptive.enabled=false;  -- 分桶时需要关闭AQE

INSERT INTO users_bucketed
SELECT * FROM raw_users;

-- 分桶表Join优化(避免Shuffle)
SELECT a.*, b.* 
FROM orders a JOIN users_bucketed b 
ON a.user_id = b.user_id;  -- 相同分桶键,避免Shuffle

2.3 DataFrame API 操作分桶

// 创建分桶表
val usersDF = spark.read.parquet("/data/raw_users")

usersDF.write
  .format("parquet")
  .bucketBy(32, "user_id")  // 32个桶,按user_id分桶
  .sortBy("registration_date")  // 可选:桶内排序
  .mode("overwrite")
  .saveAsTable("users_bucketed")

// 分桶配置
spark.conf.set("spark.sql.sources.bucketing.enabled", "true")
spark.conf.set("spark.sql.adaptive.enabled", "false")  // 分桶时临时关闭AQE

// 读取分桶表
val bucketedDF = spark.table("users_bucketed")

// 检查分桶信息
spark.sql("DESCRIBE FORMATTED users_bucketed").show()

3. 分区 vs 分桶的核心区别

3.1 对比表格

特性 分区(Partitioning) 分桶(Bucketing)
数据组织方式 按列值物理分割到不同目录 按哈希值分布到固定数量文件
目录结构 多级目录树(如/date=2024-01-01/region=north/) 固定数量的文件(如00000_0, 00001_0)
适用场景 高筛选性的查询(日期、地区等) Join操作优化、数据采样
分区/桶数 由数据决定,可能很多(数千个) 固定数量,通常2的幂次方(32,64,128)
性能优势 分区剪裁,减少I/O 避免Shuffle,优化Join
数据倾斜 可能产生倾斜(某些分区数据量大) 相对均匀分布

3.2 选择策略

def choosePartitionOrBucket(column: String, dataDF: DataFrame): String = {
    val distinctCount = dataDF.select(column).distinct().count()
    val totalCount = dataDF.count()
    
    // 分区选择标准:低基数,高筛选性
    if (distinctCount < 1000 && distinctCount.toDouble / totalCount < 0.01) {
        s"使用分区: PARTITIONED BY ($column)"
    } 
    // 分桶选择标准:高基数,常用于Join
    else if (distinctCount > 10000) {
        s"使用分桶: CLUSTERED BY ($column) INTO 64 BUCKETS"
    } else {
        "考虑组合使用分区和分桶"
    }
}

4. 分区+分桶组合使用

4.1 组合使用场景

-- 创建分区+分桶表(最佳实践)
CREATE TABLE sales_optimized (
    product_id STRING,
    amount DECIMAL(10,2),
    customer_id STRING
)
PARTITIONED BY (sale_year INT, sale_month INT)
CLUSTERED BY (product_id) INTO 64 BUCKETS
STORED AS PARQUET
LOCATION '/data/sales_optimized';

-- 插入数据
INSERT INTO sales_optimized
PARTITION (sale_year, sale_month)
SELECT 
    product_id, amount, customer_id,
    YEAR(sale_date) as sale_year,
    MONTH(sale_date) as sale_month
FROM raw_sales;

4.2 组合优势分析

graph TD
    A[原始数据] --> B[第一层: 按年月分区]
    B --> C[2024-01/]
    B --> D[2024-02/]
    B --> E[2024-03/]
    
    C --> F[第二层: 按产品ID分桶]
    F --> G[Bucket 0-63]
    F --> H[Bucket 0-63]
    
    G --> I[高效分区剪裁]
    H --> J[优化Join性能]

5. 实际应用案例

5.1 电商数据分析平台

-- 用户行为日志表设计
CREATE TABLE user_behavior (
    user_id BIGINT,
    item_id BIGINT,
    behavior_type STRING,
    timestamp BIGINT,
    province STRING,
    city STRING
)
PARTITIONED BY (dt STRING)  -- 按天分区,便于历史数据管理
CLUSTERED BY (user_id) INTO 128 BUCKETS  -- 按用户分桶,优化用户分析
STORED AS PARQUET;

-- 商品信息表设计
CREATE TABLE items (
    item_id BIGINT,
    title STRING,
    category_id BIGINT,
    price DECIMAL(10,2)
)
CLUSTERED BY (item_id) INTO 64 BUCKETS  -- 按商品分桶,优化商品分析
SORTED BY (category_id) INTO 64 BUCKETS  -- 桶内按类别排序
STORED AS PARQUET;

-- 高效查询示例
SELECT 
    u.user_id,
    i.category_id,
    COUNT(*) as behavior_count
FROM user_behavior u
JOIN items i ON u.item_id = i.item_id  -- 分桶Join,避免Shuffle
WHERE u.dt = '2024-01-15'              -- 分区剪裁
  AND u.province = 'Beijing'           -- 谓词下推
GROUP BY u.user_id, i.category_id;

5.2 性能测试对比

// 性能基准测试
def benchmarkQuery(): Unit = {
    // 测试数据
    val largeTable = spark.range(10000000).toDF("id")
    val smallTable = spark.range(1000).toDF("id")
    
    // 1. 普通表Join
    val start1 = System.currentTimeMillis()
    val result1 = largeTable.join(smallTable, "id")
    result1.count()
    val time1 = System.currentTimeMillis() - start1
    
    // 2. 分桶表Join
    largeTable.write.bucketBy(32, "id").saveAsTable("large_bucketed")
    smallTable.write.bucketBy(32, "id").saveAsTable("small_bucketed")
    
    val start2 = System.currentTimeMillis()
    val result2 = spark.table("large_bucketed").join(spark.table("small_bucketed"), "id")
    result2.count()
    val time2 = System.currentTimeMillis() - start2
    
    println(s"普通Join时间: ${time1}ms")
    println(s"分桶Join时间: ${time2}ms")
    println(s"性能提升: ${time1.toDouble / time2}x")
}

benchmarkQuery()

6. 最佳实践和注意事项

6.1 分区最佳实践

// 分区策略选择
def optimalPartitionStrategy(df: DataFrame): Seq[String] = {
    val candidateColumns = Seq("date", "region", "category")
    
    candidateColumns.filter { col =>
        val distinctCount = df.select(col).distinct().count()
        val totalCount = df.count()
        
        // 理想分区列:基数适中,查询常用
        distinctCount >= 10 && distinctCount <= 1000 &&
        df.filter(df(col).isNotNull).count() > totalCount * 0.9  // 非空值比例高
    }
}

// 避免过度分区
def validatePartitionCount(df: DataFrame, partitionCols: Seq[String]): Boolean = {
    val partitionCount = df.select(partitionCols.map(col): _*).distinct().count()
    partitionCount <= 10000  // 避免过多小文件
}

6.2 分桶最佳实践

// 分桶数选择
def calculateOptimalBuckets(df: DataFrame, bucketColumn: String): Int = {
    val dataSizeMB = df.queryExecution.optimizedPlan.stats.sizeInBytes / (1024 * 1024)
    val executorCores = spark.sparkContext.defaultParallelism
    
    // 目标:每个桶100-200MB,但不超过执行器核心数
    val targetBuckets = Math.max(1, (dataSizeMB / 150).toInt)
    Math.min(targetBuckets, executorCores * 4).nextPowerOf2  // 2的幂次方
}

// 分桶键选择标准
def goodBucketColumn(df: DataFrame, column: String): Boolean = {
    val stats = df.select(column).describe()
    val distinctCount = df.select(column).distinct().count()
    
    // 好的分桶键:高基数,分布均匀,常用于Join
    distinctCount > 1000 && !df.filter(col(column).isNull).count() > df.count() * 0.1
}

6.3 常见问题解决

// 问题1:小文件过多(过度分区)
// 解决方案:合并小文件
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")

// 问题2:分桶失效(数据倾斜)
// 解决方案:检查分桶键分布
df.select(bucketColumn).groupBy(bucketColumn).count().orderBy(desc("count")).show(10)

// 问题3:分区列选择不当
// 解决方案:分析查询模式
spark.sql("DESCRIBE HISTORY table_name").show()  // 查看查询历史

分区和分桶是Spark SQL中强大的性能优化工具。分区适合基于值的范围查询,分桶适合等值查询和Join优化。在实际应用中,通常结合使用两者以达到最佳性能。

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐