全部產品
Search
文件中心

Realtime Compute for Apache Flink:效能白皮書(Nexmark效能測試)

更新時間:May 16, 2026

本文介紹如何使用 Nexmark 基準測試載入器評測Realtime Compute Flink 版的流處理效能。

效能一覽

Nexmark 是業界通用的流處理引擎效能基準,包含 19 條標準 Query,覆蓋過濾、彙總、JOIN、視窗等典型情境。本文基於 Nexmark 測試載入器,在 8 CU 計算資源配置下,以 1 億條輸入資料為基準,對Realtime Compute Flink 版進行全量 Query 效能評測。測試結果表明:

  • 簡單查詢(如 q0、q1、q2)的 RPS 可達 400 萬~650 萬條/秒。

  • 複雜彙總與視窗查詢(如 q4、q5、q16)的 RPS 在 15 萬~63 萬條/秒之間。

整體上,Realtime Compute Flink 版的 Nexmark 效能是開源 Flink 的 3.24 倍。

測試載入器

Nexmark 是一套針對流處理引擎的標準效能基準測試。測試模型如下:

  • Nexmark 源表:按照指定 TPS 產生測試資料(Person、Auction、Bid 三類事件)。

  • Transformations:19 條標準 Nexmark Query,覆蓋過濾、轉換、彙總、JOIN、視窗等典型情境。

  • Blackhole 結果表:資料寫入 Blackhole,排除外部儲存的效能幹擾,專註評測 Flink 引擎自身的處理能力。

本文採用的 Nexmark 測試載入器基於Realtime Compute Flink 版的 OpenAPI 實現,自動化完成作業建立、部署、運行監控及結果採集的全流程。無需在控制台手動編寫 SQL 或建立作業。

測試環境

本次測試的 Flink 作業啟用了以下最佳化配置:

配置項

參數值

說明

table.exec.mini-batch.enabled

true

開啟 Mini-Batch 彙總

table.exec.mini-batch.allow-latency

2s

Mini-Batch 攢批間隔

table.optimizer.distinct-agg.split.enabled

true

開啟 Distinct 彙總拆分最佳化

execution.checkpointing.interval

3min

Checkpoint 間隔

前提條件

  • 已安裝 Java JDK 1.8.x 或更高版本。

  • 已建立工作空間,詳情請參見開通Realtime ComputeFlink版

  • 已擷取阿里雲帳號的 AccessKey ID 和 AccessKey Secret。

測試步驟

步驟一:下載測試載入器

下載 Nexmark 測試載入器壓縮包nexmark-flink.tar.gz並解壓。

解壓後的目錄結構如下:

nexmark-flink/
├── run_nexmark.sh          # 測試入口指令碼
├── nexmark_env.sh          # 環境變數設定檔(需編輯)
├── bin/                    # 運行指令碼
├── conf/                   # Flink 作業配置
├── lib/                    # JAR 包(需上傳至控制台)
└── queries-vvp/            # Nexmark Query SQL 檔案

步驟二:上傳 Nexmark JAR

  1. 登入Realtime Compute控制台

  2. 單擊進入目標專案空間,在左側導覽列檔案管理 > 上傳資源

  3. 選擇並上傳 nexmark-flink-0.2-SNAPSHOT.jar 檔案。該檔案位於測試載入器的 nexmark-flink/lib 目錄下。

  4. 上傳完成後,單擊目標檔案名稱複製 OSS 地址。該地址在後續配置參數時使用。檔案路徑格式因儲存類型而異:

    • OSS Bucket 儲存oss://<OSS Bucket 名稱>/artifacts/namespaces/<專案空間名稱>/<檔案名稱>

      例:oss://oss-test/artifacts/namespaces/flink-default/nexmark-flink-0.2-SNAPSHOT.jar
    • 全託管儲存oss://flink-fullymanaged-<工作空間ID>/artifacts/namespaces/<專案空間名稱>/<檔案名稱>

      例:oss://flink-fullymanaged-e6a123456789/artifacts/namespaces/flink-default/nexmark-flink-0.2-SNAPSHOT.jar

    如需查看工作空間的儲存類型,在Realtime Compute管理主控台單擊目標工作空間操作列下的詳情查看。

步驟三:配置運行參數

編輯 nexmark-flink/nexmark_env.sh 檔案,填寫以下參數。

參數名

說明

樣本

END_POINT

Realtime Compute Flink 版的服務存取點。根據地區選擇對應的存取點,詳情請參見服務存取點

ververica.cn-hangzhou.aliyuncs.com

AK

阿里雲帳號的 AccessKey ID。

-

SK

阿里雲帳號的 AccessKey Secret。

-

WORK_SPACE

專案工作空間的 Workspace ID。

e6a123456789

NAMESPACE

專案工作空間的 Namespace。

flink-default

NEXMARK_JAR

步驟二中上傳的 JAR 檔案的 OSS 地址。

oss://flink-fullymanaged-e6a123456789/artifacts/namespaces/flink-default/nexmark-flink-0.2-SNAPSHOT.jar

FLINK_VERSION

目標測試的 Flink 引擎版本號碼。

vvr-11.6-jdk11-flink-1.20

QUERIES

指定啟動並執行 Query。多個 Query 以逗號分隔,例如 q0q1,q2,q3。運行全部 Query 設定為 all

all

說明

運行全部 Query 耗時較長。每條 Query 需要經歷作業建立、資料產生、計算執行等階段。建議先運行單條 Query(例如將 QUERIES 設為 q0),驗證環境配置和參數填寫無誤後,再執行全量測試。

步驟四:運行測試

  1. nexmark-flink 目錄下執行以下命令。

    ./run_nexmark.sh
  2. 測試載入器通過 OpenAPI 自動建立並運行 Nexmark 作業。

  3. 運行完成後,輸出各 Query 的運行時間長度(毫秒)。樣本如下:

    INFO  com.github.nexmark.flink.vvp.Nexmark - q0 13078
    ============================================================================
    ✓ Benchmark execution completed successfully
    ============================================================================

效能表現

以下為 8 CU 計算資源配置下,開源 Flink(1.20.4)與Realtime Compute Flink 版(vvr-11.5-jdk11-flink-1.20)的 Nexmark 效能對比。每條 Query 輸入 1 億條資料,RPS = 輸入資料量 ÷ 用時。

說明

以下測試資料基於特定硬體環境和引擎版本採集。隨著底層硬體升級迭代和引擎版本更新,實際效能表現可能存在差異,測試結果僅供參考。

Query

開源 Flink on ECS

Version:1.20.4

Realtime Compute Flink 版

Version:vvr-11.5-jdk11-flink-1.20

用時(毫秒)

RPS

用時(毫秒)

RPS

RPS 相比開源 Flink(倍)

q0

58848

1,699,293

23450

4,264,392

2.51

q1

57045

1,753,002

22824

4,381,353

2.50

q2

51890

1,927,154

15224

6,568,576

3.41

q3

84986

1,176,664

21558

4,638,649

3.94

q4

553426

180,693

157117

636,468

3.52

q5

365636

273,496

357547

279,684

1.02

q7

1257452

79,526

333837

299,547

3.77

q8

79788

1,253,321

29939

3,340,125

2.67

q9

2324518

43,020

266563

375,146

8.72

q10

189985

526,357

51202

1,953,049

3.71

q11

408384

244,868

145983

685,011

2.80

q12

121554

822,680

36991

2,703,360

3.29

q14

68903

1,451,316

20012

4,997,002

3.44

q15

183709

544,339

42734

2,340,057

4.30

q16

917597

108,980

337293

296,478

2.72

q17

102847

972,318

27076

3,693,308

3.80

q18

574949

173,928

96335

1,038,044

5.97

q19

586287

170,565

95121

1,051,293

6.16

q20

1340638

74,591

231482

431,999

5.79

q21

127089

786,850

39693

2,519,336

3.20

q22

94830

1,054,519

31228

3,202,254

3.04

總計

2383209

49,695,131

9550361

15,317,480

3.24

開源 Flink 測試流程

以下為開源 Flink on ECS 的 Nexmark 測試步驟,供複現驗證。

環境準備

通過 EMR on ECS 建立 Flink 叢集,叢集配置如下:

  • EMR 版本:EMR-5.21.0

  • 硬體規格:3 台 ecs.g6a.xlarge(4 vCPU / 16 GiB),1 台 Master + 2 台 Core

  • 開啟 Hadoop 和 HDFS 服務

  • 配置好各個節點之間的免密登入。例如,將安全性群組的私密金鑰 key.pem上傳到 master 節點,然後在 master 節點的 ~/.ssh/config 寫入如下配置,注意將 IP 和檔案路徑替換為實際值:

    Host 192.168.0.0
        HostName 192.168.0.0
        User root
        IdentityFile /path/to/key.pem
        StrictHostKeyChecking no
    Host 192.168.0.1
        HostName 192.168.0.1
        User root
        IdentityFile /path/to/key.pem
        StrictHostKeyChecking no
    Host 192.168.0.2
        HostName 192.168.0.2
        User root
        IdentityFile /path/to/key.pem
        StrictHostKeyChecking no
    使用 ssh 驗證各節點間免密登入是否正常。若報錯 "bad permissions",執行 chmod 600 /path/to/key.pem 修正許可權。

軟體準備

  1. 下載目標版本的 Flink 包(Apache Flink Downloads)和 Nexmark 測試包(nexmark-flink.tgz),上傳至 Master 節點並解壓。

    tar -zxvf flink-1.20.4-bin-scala_2.12.tgz
    tar -zxvf nexmark-flink.tgz
    mv flink-1.20.4 flink
    mv nexmark-flink nexmark
  1. 將 nexmark/lib 目錄下的 JAR 檔案複製到 flink/lib,這些 JAR 包含 Nexmark 資料產生器。

    cp nexmark/lib/* flink/lib/
  1. 設定環境變數。編輯 ~/.bashrc,加入以下配置後執行 source ~/.bashrc 使其生效。

    根據實際環境配置相關路徑。
    export JAVA_HOME=/etc/alternatives/java_sdk_11
    export PATH=$JAVA_HOME/bin:$PATHexport FLINK_HOME=/mnt/disk1/flink
    export HADOOP_CLASSPATH=$(/opt/apps/HADOOP-COMMON/hadoop-common-current/bin/hadoop classpath)
  1. 配置 Flink Workers。本文使用 8 個 TaskManager,採用"Master 節點 2 個 + 每個 Core 節點各 3 個"的部署方式。

    編輯 flink/conf/workers,注意將IP替換為實際值。

    192.168.0.0
    192.168.0.0
    192.168.0.1
    192.168.0.1
    192.168.0.1
    192.168.0.2
    192.168.0.2
    192.168.0.2
  1. 用 nexmark/conf/config.yaml 替換 flink/conf/config.yaml,並更新以下配置項:

    • jobmanager.rpc.address:Master 節點 IP,如 192.168.0.0

    • state.checkpoints.dir:HDFS 路徑,如 hdfs:///checkpoints

    • taskmanager.memory.process.size4G

  1. 編輯 nexmark/conf/nexmark.yaml,將 nexmark.metric.reporter.host 設為 Master 節點 IP。

  2. 將 flink 和 nexmark 目錄及環境變數配置分發到各 Core 節點。

    注意將IP替換為實際值。
    scp -r flink 192.168.0.1:/mnt/disk1/
    scp -r flink 192.168.0.2:/mnt/disk1/
    scp -r nexmark 192.168.0.1:/mnt/disk1/
    scp -r nexmark 192.168.0.2:/mnt/disk1/
    scp ~/.bashrc 192.168.0.1:~/
    scp ~/.bashrc 192.168.0.2:~/

    分發完成後,在各 Core 節點執行 source ~/.bashrc 使環境變數生效。

  1. 在 Master 節點啟動 Flink 叢集。

    flink/bin/start-cluster.sh
  1. 初始化 Nexmark 測試環境。該指令碼會在各節點上配置 Nexmark 運行所需的 Metric Reporter。

    nexmark/bin/setup_cluster.sh

水位設定

Flink 叢集啟動後,需通過 cgroup 將各 TaskManager 進程的 CPU 水位限制在 75% 以內,避免因資源爭搶導致 TaskManager 心跳逾時失聯。

在所有運行 TaskManager 的節點上(含 Master 節點)執行以下命令:

yum install -y libcgroup libcgroup-tools
cgcreate -t root:root -a root:root -g cpu,memory:mygroup
echo 100000 > /sys/fs/cgroup/cpu/mygroup/cpu.cfs_period_us
echo 300000 > /sys/fs/cgroup/cpu/mygroup/cpu.cfs_quota_us
echo $((12 * 1024 * 1024 * 1024)) > /sys/fs/cgroup/memory/mygroup/memory.limit_in_bytes
jps | grep TaskManagerRunner | awk '{print $1}' | xargs cgclassify -g cpu,memory:mygroup

其中 cpu.cfs_quota_us / cpu.cfs_period_us = 300000 / 100000 = 3,表示該 cgroup 最多使用 3 個 CPU 核心(4 vCPU 的 75%)。

說明

Realtime Compute Flink 版(全託管環境)無需手動控制水位,可全量使用所購買的計算資源。

運行 Nexmark

在 Master 節點執行以下命令,等待運行結束後查看結果:

nexmark/bin/run_query.sh q0,q1,q2,q3,q4,q5,q7,q8,q9,q10,q11,q12,q14,q15,q16,q17,q18,q19,q20,q21,q22