
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 性能度量指标QphDSSF这是TPC-DS的官方性能指标全称是“每小时查询次数数据量比例”。计算公式相对复杂但核心思想是在特定的数据量Scale Factor, SF下系统每小时能够成功执行的查询复杂度加权几何平均数。Scale Factor (SF)数据量比例因子。SF1代表基础数据量约1GBSF1000则约1TB。你需要根据你的集群规模选择合理的SF。对于学习或小集群可以从SF1或10开始生产级评估可能需要SF100或以上。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的数据。例如生成SF100约100GB的数据并输出为CSV格式# 假设你已编译好dsdgen ./dsdgen -scale 100 -dir /path/to/output -terminate N -force Y -delimiter |这会生成24个.dat文件对应24张表字段默认以|分隔。数据上云/入湖将生成的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.partitionsShuffle阶段的分区数对性能影响巨大。通常建议设置为executor-cores * executor-instances的2-3倍。spark.sql.adaptive.enabled启用自适应查询执行AQESpark 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可能就是瓶颈。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 Testspark-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运行时间极长这几乎可以断定是数据倾斜。资源利用率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.enabledtrue默认true。Spark AQE能自动检测运行时的数据倾斜并将倾斜的分区拆分成多个小分区进行处理。这是Spark 3.x后应对倾斜的首选利器。手动处理倾斜Key分离法将倾斜的Key如NULL值或某个特定值从数据集中分离出来单独处理最后再合并结果。加盐法Salting对倾斜Key的关联字段添加随机前缀将一个大Key打散成多个小Key。这需要改写SQL逻辑。例如对于倾斜的customer_sk可以给大表和小表的这个字段都加上一个随机数后缀如0-9然后关联条件变为concat(customer_sk, ‘_’, rand(10))。这种方法比较“重”但对付极端倾斜很有效。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。验证重新运行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 SF100的负载综合能力如何能否满足未来业务增长配置建议给出一套针对当前硬件和负载推荐的Spark配置参数。后续计划指出下一步的优化方向如升级硬件、尝试Z-Ordering优化数据布局、测试Spark 3.x新特性等。6.2 建立可持续的性能基线性能测试不应是一次性的。你需要建立一个基线以便未来任何变更如Spark版本升级、集群扩容、参数调整都可以进行对比。基线版本将第一次全面测试调优前的结果和最终调优后的结果都作为基线保存下来。保存完整的结果文件spark-sql-perf生成的JSON、Spark配置文件以及测试时的监控快照。自动化脚本将数据生成、数据加载、测试执行、结果收集的流程脚本化。这样任何想进行回归测试的人都可以一键执行。变更管控任何可能影响性能的变更特别是Spark配置、集群规模、底层库版本在应用到生产环境前都应在测试环境运行一遍TPC-DS测试与基线进行对比评估影响。通过这样一套完整的实践TPC-DS就从一套陌生的SQL变成了你手中的一把精准尺子。它不仅能衡量集群的绝对性能更能帮助你深入理解Spark在应对复杂、真实负载时的行为培养出定位和解决深层性能问题的能力。下次再遇到性能问题你就不再是“感觉”快了慢了而是可以有理有据地说“根据TPC-DS的测试结果我们的集群在复杂关联查询上的QphDS100是XX其中瓶颈主要在Shuffle IO建议从以下三点进行优化……”