1. 从一次性能瓶颈排查说起:为什么我们需要TPC-DS
去年,我们团队接手了一个新的数据仓库项目,底层引擎从Hive迁移到了Spark。迁移完成后,业务方反馈说,几个核心的报表查询“感觉”变快了,但一到月底跑月度汇总大任务时,整个集群就变得异常缓慢,甚至出现过任务失败的情况。开发同学信誓旦旦地说,Spark的代码逻辑和Hive SQL是等价的,理论上性能只会更好。那问题出在哪?
我们花了一周时间,像没头苍蝇一样到处排查:是不是资源分配不合理?是不是数据倾斜了?是不是某个配置参数没调优?我们针对出问题的几个SQL做了针对性优化,确实有所改善。但老板问了一个我们答不上来的问题:“现在这套Spark集群,整体性能到底比之前的Hive强多少?它的瓶颈天花板在哪里?我们未来业务量翻倍,它还能撑得住吗?”
我们意识到,靠零散的、针对具体业务SQL的优化和感觉,无法回答这些系统性问题。我们需要一个标尺,一个能全面、客观、可重复地衡量大数据计算引擎分析能力的标尺。这就是TPC-DS进入我们视野的原因。它不是某个业务方的特定查询,而是一套公认的基准测试套件,模拟了真实决策支持系统的复杂查询负载。通过它,我们可以回答:我们的Spark集群在多并发查询、复杂关联分析、大数据量扫描等方面的真实性能如何。
所以,今天这篇内容,我就结合那次从“救火”到“建立体系”的经历,和你详细拆解如何用Spark进行TPC-DS性能测试。这不仅仅是一个跑分工具的使用教程,更是一次建立数据平台性能评估标准化的实践。你会发现,通过这个过程,你能更深刻地理解Spark的内部工作机制,并提前发现集群的潜在瓶颈。
2. TPC-DS基准测试深度解析:不止于99条SQL
很多人听到TPC-DS,第一反应就是“哦,那个有99条查询SQL的测试集”。这理解对,但太浅了。TPC-DS是一套极其严谨的工业标准,它的价值在于其高度仿真的数据模型和查询模型。
2.1 数据模型:一个高度仿真的零售业数据仓库
TPC-DS定义了一个虚构的、全球性的零售企业的数据仓库模型。它包含7张事实表和17张维度表(如store_sales销售事实表、customer客户维度表、item商品维度表等),表之间通过外键关联,形成了一个典型的星型/雪花型混合模型。这个模型涵盖了销售、库存、退货、促销、客户等多个业务主题。
更重要的是,TPC-DS的数据生成工具(dsdgen)能生成具有真实世界数据特性的测试数据:
- 数据倾斜:某些维度的数据分布不是均匀的,比如热门商品和冷门商品的销售记录量级差异巨大。
- 关联关系:表之间的关联关系(如
customer到customer_address)模拟了真实的主从关系。 - 时间序列:数据带有时间戳,支持基于时间窗口的查询。
这意味着,你的测试数据本身就不再是“理想状态”下的均匀数据,而是自带“坑点”的。测试用这种数据,结果才更有说服力。
2.2 查询负载:覆盖决策支持系统的全场景
那99条查询(Q1-Q99)和20条刷新流(RF1-RF2)是精髓所在。它们被精心设计来覆盖决策支持系统中的几乎所有操作类型:
- 简单报表类:单表聚合、过滤(如Q1)。
- 复杂关联分析:多表Join(常常是5-6张表),包括星型Join和雪花型Join(如Q19, Q42)。
- 高级分析函数:
窗口函数(Window Function)、ROLLUP、CUBE等(如Q67, Q98)。 - 迭代计算:一些查询包含了子查询或CTE(公共表表达式),形成逻辑上的迭代(如Q89)。
- 即席查询与固定报表混合:测试流程模拟了多用户并发执行查询的场景,这比单条查询顺序执行更能考验系统的并发处理能力和资源隔离能力。
刷新流(Refresh Function)模拟了数据仓库的ETL过程,在查询测试间隙执行,用于更新事实表和维度表。这考验的是Spark在处理读(查询)写(刷新)混合负载时的能力,对于流批一体或Lambda架构的评估尤为重要。
2.3 性能度量指标:QphDS@SF
这是TPC-DS的官方性能指标,全称是“每小时查询次数@数据量比例”。计算公式相对复杂,但核心思想是:在特定的数据量(Scale Factor, SF)下,系统每小时能够成功执行的查询复杂度加权几何平均数。
- Scale Factor (SF):数据量比例因子。SF=1代表基础数据量约1GB,SF=1000则约1TB。你需要根据你的集群规模选择合理的SF。对于学习或小集群,可以从SF=1或10开始;生产级评估可能需要SF=100或以上。
- Power Test:顺序执行所有查询,衡量系统处理单一流查询的原始能力。
- Throughput Test:多流并发执行查询和刷新流,衡量系统在并发压力下的吞吐能力和稳定性。
- QphDS:综合了Power Test和Throughput Test的结果计算得出,是一个综合评分。
对于我们内部评估而言,可能不需要严格计算官方的QphDS分数(那需要审计和付费)。但我们完全可以借鉴其方法论:分别进行顺序执行测试(看单任务能力)和并发执行测试(看抗压能力),并记录关键指标。
3. 搭建Spark TPC-DS测试环境:从数据生成到工具准备
工欲善其事,必先利其器。一次完整的TPC-DS测试,需要准备好数据、工具和监控。
3.1 数据生成与准备
官方数据生成工具是dsdgen。虽然Spark社区有像spark-sql-perf这样的项目可以方便地集成测试,但我建议先从dsdgen开始,理解原始数据的生成过程。
- 获取dsdgen工具:你可以从TPC官网(tpc.org)下载TPC-DS工具包,其中包含
dsdgen。或者,一些开源项目(如Apache Spark源码的tpcds-kit目录下)也提供了其源码或移植版本。 - 生成数据:在Linux环境下,使用
dsdgen生成指定SF的数据。例如,生成SF=100(约100GB)的数据,并输出为CSV格式:
这会生成24个.dat文件(对应24张表),字段默认以# 假设你已编译好dsdgen ./dsdgen -scale 100 -dir /path/to/output -terminate N -force Y -delimiter '|'|分隔。 - 数据上云/入湖:将生成的CSV数据加载到你的分布式存储系统中,如HDFS、S3或OSS。这是后续Spark读取的源头。
hadoop fs -put /path/to/output/*.dat /tpcds/sf100/
注意:直接使用CSV格式在测试中可能会因为解析开销成为性能瓶颈的一部分。对于追求极致性能的测试,建议将数据转换为列式存储格式,如Parquet或ORC。你可以用Spark先读取CSV,然后以Parquet格式写回存储。这本身也是一个对Spark IO能力的预热测试。
val df = spark.read.option("delimiter", "|").option("header", "false").schema(someSchema).csv("/tpcds/sf100/*.dat") df.write.mode("overwrite").parquet("/tpcds/sf100_parquet/")
3.2 测试工具选择:spark-sql-perf
手动组织99条SQL的执行和结果收集是噩梦。幸运的是,我们有spark-sql-perf这个开源库(由Databricks贡献)。它封装了TPC-DS的查询、数据生成以及性能结果收集。
引入依赖:如果你使用Spark Shell或自己编译项目,需要将
spark-sql-perf的jar包加入classpath。对于sbt项目,可以在build.sbt中添加:libraryDependencies += "com.databricks" %% "spark-sql-perf" % "0.5.1" // 注意版本匹配你的Spark版本对于PySpark用户,虽然核心是Scala库,但可以通过
--jars参数指定jar包路径来使用。工具核心概念:
Benchmark:基准测试类,核心对象。Tables:用于创建和管理TPC-DS表。- 通过它,你可以方便地:注册TPC-DS表结构、运行单个或全部查询、收集每条查询的执行时间、磁盘IO、CPU时间等指标。
3.3 集群与监控配置
测试环境应尽量贴近生产环境,否则测试结果没有参考价值。
Spark配置:这是性能调优的核心。你需要一个基准配置,通常从集群的默认配置开始。关键配置包括:
spark.executor.memory,spark.executor.cores,spark.executor.instances:决定了并行度。spark.sql.shuffle.partitions:Shuffle阶段的分区数,对性能影响巨大。通常建议设置为executor-cores * executor-instances的2-3倍。spark.sql.adaptive.enabled:启用自适应查询执行(AQE),Spark 3.x后强烈建议开启,它能动态优化执行计划。spark.sql.files.maxPartitionBytes:控制读取文件时每个分区的数据量,影响并行度。
我的建议是,先以一套合理的默认配置运行一遍全部查询,记录结果作为基线。然后再进行调优对比。
监控系统:性能测试不看监控,等于盲人摸象。必须配置好监控。
- Spark UI:这是最直接的。关注
Stages页面的任务执行时间、Shuffle读写量、GC时间。如果发现某个Stage时间过长,点进去看是否有数据倾斜(少数Task处理了绝大部分数据)。 - 集群监控:如YARN RM UI、Kubernetes Dashboard,查看整体CPU、内存、网络IO的使用率。确保测试期间没有其他重负载任务干扰。
- 系统级监控:如Grafana + Prometheus,监控集群节点的磁盘I/O、网络带宽。如果测试期间磁盘IO持续100%,那么IO可能就是瓶颈。
- Spark UI:这是最直接的。关注
4. 执行测试与核心性能指标采集
环境准备好后,就可以开始跑测试了。这个过程是自动化的,但需要你仔细观察。
4.1 编写测试脚本
以下是一个使用spark-sql-perf和Scala的示例脚本框架:
import com.databricks.spark.sql.perf.tpcds.TPCDS import org.apache.spark.sql.SparkSession val spark = SparkSession.builder() .appName("TPC-DS Benchmark") .config("spark.sql.adaptive.enabled", "true") .config("spark.sql.shuffle.partitions", "200") // ... 其他配置 .getOrCreate() // 1. 创建TPCDS实例 val tpcds = new TPCDS(spark.sqlContext) // 2. 指定数据位置和格式(例如Parquet) val databaseName = "tpcds_sf100" val dataLocation = "/tpcds/sf100_parquet" val tables = tpcds.tables tables.createExternalTables(dataLocation, "parquet", databaseName, overwrite = true, discoverPartitions = false) // 3. 设置实验 import com.databricks.spark.sql.perf._ val experiment = tpcds.runExperiment( queries = tpcds.queries, // 运行所有查询 iterations = 1, // 迭代次数,多次运行取平均可以消除偶然性 resultLocation = "/tmp/tpcds_results", // 结果保存路径 tags = Map("runType" -> "power", "scaleFactor" -> "100") // 打标签 ) // 4. 等待实验完成,并生成报告 experiment.waitForFinish() val result = experiment.getCurrentResults() result.show(false) // 查看详细结果 // 5. 可以将结果进一步分析或保存 spark.read.json("/tmp/tpcds_results/*.json").createOrReplaceTempView("results")对于吞吐量测试(Throughput Test),spark-sql-perf也提供了并发执行的支持,你需要定义多个Query流并同时运行。
4.2 关键性能指标收集
运行过程中,除了最终的总耗时,更要关注以下微观指标,它们能告诉你瓶颈在哪:
- 查询执行时间:每条SQL从开始到结束的时间。这是最直观的指标。列出最慢的10条查询,它们是重点分析对象。
- Stage执行时间与Shuffle数据量:在Spark UI中,查看耗时最长的Stage。重点关注:
- Shuffle Read/Write Size:过大的Shuffle数据量会拖慢速度,并消耗大量磁盘和网络IO。这可能意味着
spark.sql.shuffle.partitions设置不合理,或者发生了数据倾斜。 - GC Time:如果GC时间占比很高(比如超过10%),说明Executor内存配置或使用方式可能有问题,频繁Full GC会严重拖慢任务。
- Task Duration分布:如果某个Stage里大部分Task很快完成,但少数几个Task运行时间极长,这几乎可以断定是数据倾斜。
- Shuffle Read/Write Size:过大的Shuffle数据量会拖慢速度,并消耗大量磁盘和网络IO。这可能意味着
- 资源利用率:
- CPU利用率:在整个测试期间,集群的CPU使用率是否平稳且较高(如70%以上)?如果CPU使用率很低,但任务很慢,可能是IO瓶颈或任务并发度不够。
- 内存利用率:Executor内存是否充足?是否频繁发生Spill(溢写)到磁盘?Spill会极大降低性能。
- 磁盘IO和网络IO:监控这些指标,看是否在测试期间达到瓶颈。
4.3 常见问题与初步分析
第一次跑TPC-DS,你大概率会遇到以下问题:
- OOM(内存溢出):最常见。可能原因:1)
spark.executor.memory设置过小;2) 发生了严重的数据倾斜,导致单个Task处理的数据量远超预期;3) 查询本身(如CUBE)会产生巨大的中间结果集。 - 个别查询极慢:去Spark UI分析该查询的执行计划。重点关注:
- 是否出现了
Cartesian Product(笛卡尔积)?这通常是性能杀手。 - Join的先后顺序是否合理?大表是否最后才参与Join?
- 是否可以用
广播连接(Broadcast Join)优化小表的Join?检查spark.sql.autoBroadcastJoinThreshold设置。
- 是否出现了
- 并发测试时任务排队严重:说明集群资源不足,或者
spark.dynamicAllocation配置需要调整,确保有足够的Executor来应对并发查询。
5. 基于测试结果的深度调优实战
拿到基线测试结果后,调优才真正开始。调优不是盲目改参数,而是基于证据的针对性优化。
5.1 针对数据倾斜的优化
这是TPC-DS测试中最常遇到,也是收益最明显的优化点。
如何识别:在Spark UI的Stage详情页,查看任务执行时间的“任务时间线”或“任务数据量分布”。如果发现少数几个任务的输入数据量(Input Size)或Shuffle数据量是其他任务的几十倍甚至上百倍,那就是倾斜。
优化手段:
- 增加Shuffle分区数:通过
spark.sql.shuffle.partitions(默认200)增加分区数,让数据被打散到更多的Task中处理。这对于轻度倾斜可能有效。 - 使用AQE的倾斜Join优化:确保
spark.sql.adaptive.skewJoin.enabled=true(默认true)。Spark AQE能自动检测运行时的数据倾斜,并将倾斜的分区拆分成多个小分区进行处理。这是Spark 3.x后应对倾斜的首选利器。 - 手动处理倾斜Key:
- 分离法:将倾斜的Key(如
NULL值或某个特定值)从数据集中分离出来,单独处理,最后再合并结果。 - 加盐法(Salting):对倾斜Key的关联字段添加随机前缀,将一个大Key打散成多个小Key。这需要改写SQL逻辑。例如,对于倾斜的
customer_sk,可以给大表和小表的这个字段都加上一个随机数后缀(如0-9),然后关联条件变为concat(customer_sk, ‘_’, rand(10))。这种方法比较“重”,但对付极端倾斜很有效。
- 分离法:将倾斜的Key(如
5.2 SQL与执行计划优化
Spark SQL的Catalyst优化器很强大,但并非万能。我们需要理解其执行计划。
- 查看与分析执行计划:使用
df.explain(true)或Spark UI的SQL页。关注== Physical Plan ==部分。 - 避免隐式类型转换:确保Join或Where条件两边的数据类型一致,否则会引发昂贵的类型转换,且可能阻止索引(如果有的话)或Bloom Filter的使用。
- 谓词下推:确保过滤条件(Where)能尽可能早地在数据扫描阶段就执行,减少流入下游的数据量。使用Parquet/ORC格式时,此优化通常是自动的。
- Join策略选择:Spark有多种Join策略(BroadcastHashJoin, SortMergeJoin, ShuffleHashJoin)。通常,小表(小于
spark.sql.autoBroadcastJoinThreshold,默认10MB)会自动进行广播。对于TPC-DS,有些维度表可能因为数据生成方式而“显得”不大,但关联后数据量很大,需要观察是否选择了正确的Join策略。有时可以手动通过/*+ BROADCAST(t) */提示来强制广播。
5.3 集群与资源配置调优
这是硬件和框架层面的调优。
- Executor配置黄金法则:
- 每个Executor的核数:建议在3-5个之间。太少不利于并行利用多核,太多会导致Executor内任务竞争资源,且GC压力大。我通常从4开始。
- Executor内存:根据每个核的内存来定。通常每个核分配4-8GB内存是一个不错的起点。例如,4核Executor可以配16GB-32GB内存。
- 内存结构:
spark.executor.memory分为Storage、Execution和Reserved。通过spark.memory.fraction(默认0.6)和spark.memory.storageFraction(默认0.5)来调整。如果任务缓存需求大,可以适当提高Storage部分;如果Shuffle多,可以适当提高Execution部分。
- 动态分配:对于并发吞吐测试,开启
spark.dynamicAllocation.enabled可以让Spark根据任务队列动态申请和释放Executor,提高资源利用率。 - Shuffle服务与I/O:如果使用YARN,确保启用
spark.shuffle.service.enabled,这样Executor释放后,Shuffle数据不会丢失。使用高性能本地磁盘(如SSD)作为spark.local.dir,可以显著提升Shuffle和Spill的性能。
5.4 一个调优迭代案例
假设我们发现Q72查询(一个多表Join的复杂查询)在基线测试中特别慢。
- 定位:在Spark UI中找到Q72的Job,发现其中一个Stage的Shuffle Write高达500GB,且有一个Task处理了其中300GB的数据(严重倾斜)。
- 分析:检查该Stage的执行计划,发现是一个大表
store_sales和date_dim的Join,关联键ss_sold_date_sk在date_dim表中分布均匀,但在store_sales表中,某些日期(如促销日)的销售记录异常多。 - 优化:
- 尝试一(AQE):确认
spark.sql.adaptive.skewJoin.enabled已开启。重新运行,观察AQE是否自动拆分了倾斜分区。如果有效,时间会大幅下降。 - 尝试二(手动加盐):如果AQE效果不明显(例如倾斜过于极端),考虑手动优化。改写Q72的SQL,对
store_sales表的ss_sold_date_sk在特定日期范围内添加随机前缀,并对date_dim表做相应扩展。这一步需要深入理解SQL逻辑,改动复杂。 - 尝试三(调整资源):同时,我们可能发现该查询的Executor内存不足,导致频繁Full GC。将
spark.executor.memory从16G提升到24G,并增加spark.sql.shuffle.partitions到400。
- 尝试一(AQE):确认
- 验证:重新运行Q72,并对比优化前后的Stage时间、Shuffle数据量和GC时间。记录优化效果。
6. 测试报告撰写与性能基线建立
测试和调优的最终目的,是形成一份有价值的报告,并建立一个可持续比较的性能基线。
6.1 测试报告的核心要素
一份好的内部测试报告,不应只是一堆数字,而应有分析、有结论、有建议。
- 测试概述:说明测试目的、集群硬件配置(节点数、CPU、内存、磁盘类型、网络)、软件版本(Spark、Hadoop、JDK)、测试数据量(SF)、测试类型(Power/Throughput)。
- 详细结果:
- 总览:总耗时、成功查询数、失败查询数。
- 查询性能分布:以柱状图或百分位表(P50, P90, P99)展示所有查询的执行时间分布。列出Top 10最慢查询及其耗时。
- 资源使用情况:测试期间集群平均/峰值CPU、内存、磁盘IO、网络IO使用率图表。
- 关键事件:记录测试过程中出现的OOM、长GC、数据倾斜等异常情况。
- 深度分析:
- 瓶颈分析:结合Spark UI和监控数据,指出系统瓶颈是CPU、内存、磁盘IO还是网络?是某个特定类型的查询(如多表Join)还是所有查询都慢?
- 对比分析:(如果有)与旧系统(如Hive)或不同配置/版本的Spark集群进行对比,用数据说明性能提升或下降的百分比。
- 调优总结:列出本次测试中实施的有效调优手段(如调整了哪个参数,从多少调到多少,效果如何),以及无效的尝试。
- 结论与建议:
- 性能结论:当前集群处理TPC-DS SF=100的负载,综合能力如何?能否满足未来业务增长?
- 配置建议:给出一套针对当前硬件和负载推荐的Spark配置参数。
- 后续计划:指出下一步的优化方向(如升级硬件、尝试Z-Ordering优化数据布局、测试Spark 3.x新特性等)。
6.2 建立可持续的性能基线
性能测试不应是一次性的。你需要建立一个基线,以便未来任何变更(如Spark版本升级、集群扩容、参数调整)都可以进行对比。
- 基线版本:将第一次全面测试(调优前)的结果和最终调优后的结果,都作为基线保存下来。保存完整的结果文件(
spark-sql-perf生成的JSON)、Spark配置文件以及测试时的监控快照。 - 自动化脚本:将数据生成、数据加载、测试执行、结果收集的流程脚本化。这样,任何想进行回归测试的人,都可以一键执行。
- 变更管控:任何可能影响性能的变更(特别是Spark配置、集群规模、底层库版本),在应用到生产环境前,都应在测试环境运行一遍TPC-DS测试,与基线进行对比,评估影响。
通过这样一套完整的实践,TPC-DS就从一套陌生的SQL,变成了你手中的一把精准尺子。它不仅能衡量集群的绝对性能,更能帮助你深入理解Spark在应对复杂、真实负载时的行为,培养出定位和解决深层性能问题的能力。下次再遇到性能问题,你就不再是“感觉”快了慢了,而是可以有理有据地说:“根据TPC-DS的测试结果,我们的集群在复杂关联查询上的QphDS@100是XX,其中瓶颈主要在Shuffle IO,建议从以下三点进行优化……”