- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
本指南以 SeaTunnel 仓库中的docs/zh/engines/zeta/benchmark.md为骨架,系统讲解 Zeta 引擎基准测试的动机、JMH 测量架构、指标含义、本地构建与运行方法、PR 性能对比与 Profiler 诊断流程。读完本文,你将能够自行构建 Benchmark Runner、用 JMH JSON 生成报告、在本地跑通 CPU/锁/GC 诊断,并正确解读 Baseline 与 Candidate 的对比结果,判断一次改动是否真的带来了性能改善。
为什么需要基准测试
基准测试的目标,是让 SeaTunnel Zeta 引擎在持续演进中运行得更稳定、处理得更快,并更高效地利用计算与存储资源。
随着数据规模增长和使用场景丰富,引擎需要在更高负载下保持吞吐、控制延迟,并承担 Checkpoint、状态存储(IMap)和可观测性等能力的开销。基准测试帮助我们发现限制处理能力的瓶颈,识别负载增长时出现的性能波动,为提升引擎的处理效率和运行稳定性提供依据。
它也为社区提供了一套共同的性能验证方式:让优化收益可以被量化,让潜在的性能回退更早被发现,让不同贡献者能够复现和比较结果。通过持续积累基线与测试场景,我们可以更有依据地评估每次改动,推动引擎性能持续改进。
基准测试架构
SeaTunnel 的基准测试建立在 JMH(Java Microbenchmark Harness)之上,被测对象是嵌入式运行的 Zeta 引擎。整体结构可以概括为:JMH 负责测量纪律,SeaTunnel Client 负责提交作业,嵌入式 Zeta Cluster 负责执行真实的数据处理流程。
从源码看,这一结构在 SeaTunnelEnvironmentContext.java 中有直接对应实现:它是一个@State(Scope.Thread)的 JMH 状态对象,owns an embedded, single-node Zeta cluster and its client——集群在每次 JMH Trial 启动时创建一次,而每个 Benchmark invocation 提交并等待一个完整的有限(bounded)作业完成。类中定义了SLOT_COUNT = 12、SOURCE_START_DELAY_MILLIS = 250L、SOURCE_EMIT_BATCH_SIZE = 1_024等常量,并加载source-sink.conf.template、source-transform-sink.conf.template、engine.yaml.template等作业模板,保证每个被测流水线的负载条件一致。
职责与生命周期
- JMH负责独立 JVM(Fork)、预热(Warmup)、测量(Measurement)和结果采集;
- Environment Context负责准备测试数据,并在需要时启动嵌入式 Zeta 运行环境;
- 被测方法(
@Benchmark方法)调用要研究的生产操作,例如提交一条 source→sink 流水线并等待完成; - 测试结束后统一释放资源(TearDown)。
测试数据准备、环境启动和清理通常放在计时之外;只有它们本身就是研究对象时,才纳入测量。基准测试应明确这个边界,并校验被测工作确实产生了有效结果。
运行配置
共享 JMH 配置在 BenchmarkBase.java 中定义,与文档描述完全一致:
@State(Scope.Thread) @OutputTimeUnit(TimeUnit.MILLISECONDS) @BenchmarkMode(Mode.Throughput) @Fork(3) @Warmup(iterations = 3) @Measurement(iterations = 5) public abstract class BenchmarkBase {}即共享默认值为3 个 fork、3 次预热、5 次测量。每个 fork 在独立 JVM 中运行,预热在用于计算 Score 的样本采集之前完成;方法注解和命令行参数可以覆盖共享默认值。
线程数、堆大小、垃圾回收器和 JVM 可见处理器数都属于实验条件,应记录最终生效的值,并在跨版本对比时保持一致。注意:限制 JVM 可见处理器数量(如-XX:ActiveProcessorCount)不等于操作系统级 CPU 绑核。以流水线基准 SeaTunnelPipelineBenchmark.java 为例,它通过@Fork的jvmArgsAppend固定了实验条件:
@Fork(value = 3, jvmArgsAppend = { "-Xms4g", "-Xmx4g", "-XX:+UseG1GC", "-XX:+AlwaysPreTouch", "-XX:+DisableExplicitGC", "-XX:ActiveProcessorCount=4", "-Djava.net.preferIPv4Stack=true" })同时该方法类通过@Param暴露负载参数,默认值为offeredRatePerSecond=600000(每秒供给速率)、parallelism=4、payloadSize=256(每行字符数)、transformOperations=64(transform 操作数),每次 invocation 处理RECORDS_PER_INVOCATION = 1_000_000条记录(见@OperationsPerInvocation)。这些参数都可以用 JMH 的-p在命令行覆盖。
环境与可复现性
资源需求取决于所选负载。基准测试建议:
- 尽量减少机器上的无关任务;
- 在同一台机器上运行 Baseline 和 Candidate;
- 随原始结果记录JDK、JVM 设置、输入参数和代码版本(Git SHA)。
基准测试的结论适用于它所测量的操作和负载。评估更广泛的生产收益时,还需要在对应的部署条件下验证。
指标与结果解读
JMH 指标
一行 JMH 结果由“测什么、怎么测、结果多少、结果有多稳定”四部分组成:
Benchmark (Parameters) Mode Cnt Score Error Units| 字段 | 含义 | 解读方式 |
|---|---|---|
Benchmark/Parameters | 被测方法与负载参数 | 对比时必须一致。 |
Mode | thrpt测吞吐;avgt测平均耗时;sample测耗时分布;ss测单次执行耗时 | thrpt越大越好,其余耗时模式越小越好。 |
Cnt | 参与统计的测量样本数,不包含预热 | 吞吐与平均耗时模式通常为fork 数 × 测量 iteration 数。 |
Score | 测量样本的平均性能 | Score = Σxᵢ / n。结合 Mode 和 Units 判断方向。 |
Error | Score 置信区间的半宽,与 Score 单位相同 | 置信区间 = [Score − Error, Score + Error]。越小表示均值估计越精确。 |
Units | Score 的单位 | ops/time表示吞吐;time/op表示单次操作耗时。 |
CV | SeaTunnel 报告计算的样本相对波动 | CV = 样本标准差 / abs(Score) × 100%。越小表示样本越集中。 |
其中xᵢ为单个测量样本,n为 Cnt,abs表示绝对值。比较不同结果前,还要确认 JDK、线程数和 JVM 参数一致。
对比报告指标
B为 Baseline,C为 Candidate;median为各轮有效结果的中位数。B/C 分别按各自版本计算,比较前先核对 SHA、方法、参数和运行环境。
| 字段 | 用途 | 计算方式 | 如何解读 |
|---|---|---|---|
| Benchmark | 标识被测方法 | 按完整方法名与参数配对 | 表中显示简写名称。 |
| Parameters | 标识负载条件 | 取测试参数 | 两侧应一致。 |
| Score B / C | 两个版本的代表性能 | median(各轮 Score) | 吞吐越大越好,耗时越小越好。 |
| Score Change | 性能变化幅度 | 吞吐:(C / B − 1) × 100%;耗时:(1 − C / B) × 100% | 正值改善,负值回退。 |
| CV B / C | 两个版本的样本波动 | median(各轮 CV) | 越小表示样本越集中。 |
| CV Change | 波动变化幅度 | (CV C / CV B − 1) × 100% | 负值表示波动减小。 |
| Error B / C | 两个版本的相对不确定性 | median(各轮 Error / abs(Score) × 100%) | 百分比,不是原始 JMH 的绝对 Error。 |
| Error Change | 相对不确定性变化幅度 | (Error C / Error B − 1) × 100% | 负值表示相对不确定性减小。 |
| Unit | Score 的共同单位 | 如ops/ms、us/op | 其余数值列均为百分比。 |
abs表示绝对值。缺少有效数据或变化公式的分母为 0 时显示n/a;0.00%可能来自舍入,复算使用原始 JSON。
这些计算逻辑在 regression_report.py 中有完整实现:Score 取各轮中位数并做单位换算(ops/s、ops/ms、ops/us、ops/ns之间按比例折算),Score Change 按direction字段做方向调整(耗时类指标lower方向取反),CV 为样本标准差除以均值,Error 为相对误差。
变化幅度与统计结论:表中的变化幅度不是显著性检验;Workflow 执行成功也不等于性能没有回退。
Pipeline 级结果指标
除了纯 JMH 指标,SeaTunnel 还针对完整的嵌入式流水线运行输出 Pipeline 结果(由 SeaTunnelEnvironmentContext.java 采集、save_jmh_result.py归一化、regression_report.py渲染),常见列包括:
- Throughput:以
rows/s计的行吞吐; - P50 / P95 / P99 / Max:事件时间(event-time)延迟的各分位数与最大值,单位为
ms;超出可测范围的百分位会以“下界”形式显示,例如>60,000; - Growth:无单位的延迟增长比率,计算方式为
(后一半 P99 + 1) / (前一半 P99 + 1),用于观察长时间运行下的延迟漂移; - Valid(C/S/O):只统计测量样本(不含预热),
✅ x/y表示全部样本完整(complete)、可持续(sustainable)且无延迟溢出(overflow);否则以C/S/O分别给出三者的样本数。
判断结果是否可信
先确认运行的是目标版本与方法、输出通过了正确性检查,且两个版本使用相同设置。再结合 Score、Error、CV 以及各个 fork/iteration 的原始样本判断变化。
| 观察结果 | 解读与下一步 |
|---|---|
| 多轮运行都表现出一致改善,且波动较小 | 连同负载条件和测量边界一起报告收益。 |
| 差异接近波动幅度,或多轮变化方向不一致 | 暂不下结论,在受控环境下复测。 |
| 同一版本的多数方法同时明显变化 | 先检查机器负载、CPU 频率、JDK 和环境信息,再判断是否来自代码。 |
| 某个方法持续回退 | 使用 Profiler 定位新增开销,再重复不带 Profiler 的对比。 |
保留全部样本。单次更好的 iteration 或较大的提升百分比,都不足以单独证明改善可以重复。
可视化
使用-rf json -rff <file>生成 JMH JSON,可导入 JMH Visualizer 网页工具,按方法名和参数比较 Score、Error、fork 和 iteration。图表可能将多个参数值组合成标签,应结合图例和原始 JSON 确认各组实验条件。分享图表时保留原始 JMH 文件,让其他人能够查看底层样本。
仓库中的两个 Python 工具可以生成标准化产物:
- save_jmh_result.py:把原始 JMH JSON 与 Pipeline 结果归一化为标准 JSON 报告;
- regression_report.py:基于标准报告渲染 JMH 对比表与 Pipeline 对比表的 Markdown 摘要。
本地运行
构建与准备
在仓库根目录执行命令,启用benchmarkprofile,构建 JMH Runner:
./mvnw -Pbenchmark -pl seatunnel-benchmarks -am -DskipTests package git rev-parse HEAD java -versionRunner 产物为seatunnel-benchmarks/target/benchmarks.jar。切换代码版本或修改 Benchmark 后必须重新构建:当前 Git HEAD 不能证明已有 JAR 包含该版本。未提交的生产代码或 Fixture 改动也应与 SHA 一起记录。
在 IntelliJ IDEA 中,启用 MavenProfiles下的benchmark,点击Reload All Maven Projects。如果仍未显示模块,将seatunnel-benchmarks/pom.xml添加为 Maven 项目后重新加载。
通过 JAR 执行
通过-l列出可用方法。以下示例中的<benchmark-method>应替换为输出中的完整方法名,保留末尾的$;再用-lp查看它支持的参数:
java -jar seatunnel-benchmarks/target/benchmarks.jar -l java -jar seatunnel-benchmarks/target/benchmarks.jar \ '<benchmark-method>$' -lp使用方法配置的预热、测量和 fork 运行,并保存 JMH JSON:
java -jar seatunnel-benchmarks/target/benchmarks.jar \ '<benchmark-method>$' \ -rf json -rff seatunnel-benchmarks/target/benchmark-result.json选择器是正则表达式。使用完整方法名并在末尾加$,即可选择一个方法。长时间运行前先用-l确认匹配范围。
短跑只用于功能验证:冒烟验证可以追加
-f 1 -wi 1 -i 1 -w 1s -r 1s,短跑结果仅用于确认功能可用。
使用-p覆盖所选方法支持的负载参数。将<parameter>替换为-lp列出的参数名,将<value>替换为要测试的值:
java -jar seatunnel-benchmarks/target/benchmarks.jar \ '<benchmark-method>$' \ -p '<parameter>=<value>' \ -rf json -rff seatunnel-benchmarks/target/benchmark-result.json研究负载的影响时,每轮只改变一个参数。跨版本对比时,选择器、参数、线程数、JDK 和 JVM 设置应保持一致。
内置的 Benchmark 方法与测试套件
从 SeaTunnelPipelineBenchmark.java 可以看到,流水线类提供了 5 个被测方法,用于量化不同能力叠加对吞吐与延迟的影响:
sourceSink:纯 source→sink 数据通路;sourceTransformSink:叠加 transform 处理;sourceTransformSinkWithObservability:叠加可观测性指标采集;sourceTransformSinkWithTrace:叠加链路追踪;sourceTransformSinkWithObservabilityAndTrace:可观测性与追踪同时开启。
仓库还预置了核心测试套件 benchmarks_core.txt,覆盖基础数据通路(SeaTunnelRowBenchmark、IntermediateQueueBenchmark、DebeziumJsonFormatBenchmark、SeaTunnelPipelineBenchmark.sourceSink$、SeaTunnelPipelineBenchmark.sourceTransformSink$)、Checkpoint 协调与存储(CheckpointingTimeBenchmark.checkpointSingleInput$、CheckpointStorageBenchmark.checkpointPersistenceTransaction$)、高频 IMap 状态路径(IMapJobStorageBenchmark.taskGroupStateTransition$、runningMetricsReport$)以及 DAG 持久化与重载(IMapDagStorageBenchmark.finishedJobDagStore$、finishedJobDagLoad$)。CI 中的benchmarks输入项可以选择这些预设套件,也可以直接用custom_benchmarks填写精确方法选择器。
性能诊断与对比
PR 对比
BenchmarksWorkflow 用相同的负载和运行环境比较 Baseline 与 PR,回答一个核心问题:这次改动让目标操作变快了,还是引入了回退?
Workflow 参数
进入 GitHub Actions,选择Benchmarks,点击Run workflow:
| 输入项 | 填写内容 |
|---|---|
Use workflow from | Workflow 所在分支,通常选择dev;它不代表被测代码版本。 |
seatunnel_ref | Baseline 的分支、Tag 或 SHA;推荐填写固定 SHA。 |
pr_number | Candidate PR 的数字编号;留空则只测试 Baseline。 |
benchmarks | 选择预设的测试套件或测试项。 |
custom_benchmarks | 可选,填写精确方法<benchmark-method>$;填写后覆盖benchmarks。 |
Baseline 和 Candidate 必须包含相同的测试方法与 Fixture,否则结果无法配对。不要使用.*做日常 PR 对比,它会运行全部方法和参数组合;只选择改动影响的操作即可。
Workflow 会在同一个 Worker 上按以下顺序交替运行,降低机器状态随时间变化带来的偏差:
Baseline → Candidate → Candidate → Baseline这一ABBA 交替顺序在 run_benchmarks.sh 中直接实现:有pr_number时依次执行baseline-1 → candidate-1 → candidate-2 → baseline-2,每轮只使用 1 个 fork(注释说明 ABBA 序列本身已提供每个版本两次独立的 fork JVM 运行,外层再各用 1 个 fork 可避免默认对比超过作业时限);无 PR 时仅运行一次 Baseline(3 个 fork),并直接生成单版本报告。脚本还会在运行前把uname -a、lscpu、nproc、free -h、java -version等信息写入environment.txt,并抓取 CPU 型号与内存总量作为报告的环境元数据。
Java 8 和 Java 11 分别执行这组对比,报告汇总每个版本的两轮结果。脚本对.*全量选择会给出警告:完整套件对比可能超过 240 分钟的 Workflow 时限,应优先使用benchmarks_core、某个 Benchmark 类或自定义选择器。
运行完成后,确认 Job 已执行到 JMH 测量,并核对 Summary 中的 SHA、方法和参数。各字段及变化公式见上文对比报告指标;结论不明确时应重新运行。原始 JMH JSON、标准化报告和环境信息可从 artifact 下载。
性能诊断
诊断用于解释性能变化:使用 Profiling 解释已经观察到的性能变化、异常 Score 或较高的 Error/CV。Profiler 会引入额外开销,因此诊断报告与正常报告完全分开,诊断 Score 不能用于性能回归比较。
诊断选择器必须且只能匹配一个 benchmark 方法;.*或能够匹配多个方法的类名会被拒绝。
Workflow 参数
Benchmarks Diagnostics每次诊断一个版本,不执行 Baseline/Candidate 对比。
| 输入项 | 填写方式 |
|---|---|
Use workflow from | 诊断 Workflow 与工具所在分支,通常为dev。 |
seatunnel_ref | 未选择 PR 时,要诊断的分支、Tag 或 SHA。 |
pr_number | 可选可信 PR 编号;填写后以 PR Head 替代seatunnel_ref作为诊断目标。 |
benchmark | 填写从-l查询到的方法选择器<benchmark-method>$;匹配多个方法的选择器会被拒绝。 |
java_version | 8或11,与待分析的正常运行保持一致。 |
profile | cpu查看执行热点,wall查看含等待在内的耗时栈,lock查看锁竞争,gc查看分配与 GC 指标;all分别运行这四种模式。 |
capture_jfr | 勾选后增加一次独立 JFR 录制,用于离线分析。 |
jmh_args | 可选 JMH 参数或负载参数,以所选方法的-lp输出为准;留空使用默认值,fork 固定为 1。 |
本地运行 Profiler
使用同一脚本可以在本地诊断已构建的 Benchmark,对应实现为 profile_benchmarks.sh。下面以 CPU 分析为例;将cpu替换为wall、lock或gc即可切换模式:
bash tools/benchmarks/profile_benchmarks.sh profile cpu \ --benchmark '<benchmark-method>$' bash tools/benchmarks/profile_benchmarks.sh capture jfr --benchmark '<benchmark-method>$'| 参数 | 用途 |
|---|---|
profile <mode> | 选择 CPU 热点、耗时栈、锁竞争或 GC 分配分析。 |
--benchmark | 指定一个精确的 Benchmark 方法。 |
--repository | 可选,指定已构建 Benchmark JAR 的代码目录。 |
--output | 可选,指定一个不存在或内容为空的输出目录。 |
-- <JMH 参数> | 可选,覆盖预热、测量或负载参数。 |
几个值得注意的实现细节(均可在脚本中核实):
- CPU、wall-clock 和 lock 模式需要安装async-profiler并设置
ASYNC_PROFILER_HOME(脚本会查找$ASYNC_PROFILER_HOME/lib/libasyncProfiler.so或.dylib,并用jfrconv把录制结果转换成火焰图);GC 和 JFR 使用 JMH 内置 Profiler(gc:alloc=true;churn=true;churnWait=500、jfr:configName=profile;stackDepth=256); - 诊断运行固定使用一个 fork:脚本会检查
-f参数,未指定时自动追加-f 1,指定了非 1 的值会直接报错退出; - 运行前脚本会先用
-l解析选择器,匹配数量不是恰好 1 个时拒绝执行; - 默认在
seatunnel-benchmarks/target/profiles下按命令-模式-时间戳创建独立输出目录; - lock 模式以 10000 纳秒作为锁竞争持续时间阈值(低阈值能捕获短期的 monitor 竞争),CPU 事件默认取
cpu且可通过ASYNC_PROFILER_CPU_EVENT覆盖。
查看诊断产物
先在 Job Summary 中确认目标版本、测试参数和采样结果,再下载对应模式的 artifact。CPU、wall-clock 和 lock 模式提供火焰图(脚本同时生成正向与反向两张 HTML),GC 模式提供分配与回收摘要;JMH 日志和 JSON 用于复核本次运行。启用capture_jfr时还会生成可供离线分析的 JFR 文件。
lock 模式显示 0 个样本通常表示本次运行未观察到锁竞争。Profiler 会改变程序执行成本,因此诊断结果只用于定位原因,性能提升或回退仍应由不带 Profiler 的 PR 对比确认。
参考论文
SeaTunnel 的基准测试方法论参考了以下关于严谨 Java 性能评估与分布式流处理基准的经典研究:
- Andy Georges、Dries Buytaert、Lieven Eeckhout,Statistically Rigorous Java Performance Evaluation,OOPSLA 2007。
- Tomas Kalibera、Richard Jones,Rigorous Benchmarking in Reasonable Time,ISMM 2013。
- Jeyhun Karimov 等,Benchmarking Distributed Stream Data Processing Systems,ICDE 2018。
前两篇为 JMH 场景下的统计严谨性(置信区间、预热与样本量)提供了方法论基础;第三篇针对分布式流处理系统,与本仓库中嵌入式集群、吞吐/延迟双维指标的测试设计相呼应。
- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
相关推荐
SeaTunnel Zeta 引擎基准测试(Benchmark)完全指南:架构、指标解读与本地执行
SeaTunnel Zeta 引擎基准测试(Benchmark)完全指南:架构、指标解读与本地执行 SeaTunnel 的 Zeta 引擎在数据量增长和使用场景
数据集成ETL大数据批处理流处理变更数据捕获颠覆传统:3行命令实现抖音无水印视频全链路管理的开源方案
颠覆传统:3行命令实现抖音无水印视频全链路管理的开源方案 douyin downloader是一款专注于抖音内容高效获取的开源工具,通过智能解析中枢与多线程引擎
构建工具CLISeaTunnel Engine(Zeta)引擎全解析:SeaTunnel 原生执行引擎的架构设计与实战入门
SeaTunnel Engine(Zeta)引擎全解析:SeaTunnel 原生执行引擎的架构设计与实战入门 SeaTunnel Engine 是 SeaTun
数据集成ETL大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考