本文介紹如何使用 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
單擊進入目標專案空間,在左側導覽列。
選擇並上傳
nexmark-flink-0.2-SNAPSHOT.jar檔案。該檔案位於測試載入器的nexmark-flink/lib目錄下。上傳完成後,單擊目標檔案名稱複製 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 以逗號分隔,例如 | all |
運行全部 Query 耗時較長。每條 Query 需要經歷作業建立、資料產生、計算執行等階段。建議先運行單條 Query(例如將 QUERIES 設為 q0),驗證環境配置和參數填寫無誤後,再執行全量測試。
步驟四:運行測試
在
nexmark-flink目錄下執行以下命令。./run_nexmark.sh測試載入器通過 OpenAPI 自動建立並運行 Nexmark 作業。
運行完成後,輸出各 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修正許可權。
軟體準備
下載目標版本的 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
將
nexmark/lib目錄下的 JAR 檔案複製到flink/lib,這些 JAR 包含 Nexmark 資料產生器。cp nexmark/lib/* flink/lib/
設定環境變數。編輯
~/.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)
配置並啟動 Flink 叢集
配置 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
用
nexmark/conf/config.yaml替換flink/conf/config.yaml,並更新以下配置項:jobmanager.rpc.address:Master 節點 IP,如192.168.0.0state.checkpoints.dir:HDFS 路徑,如hdfs:///checkpointstaskmanager.memory.process.size:4G
編輯
nexmark/conf/nexmark.yaml,將nexmark.metric.reporter.host設為 Master 節點 IP。將
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使環境變數生效。
在 Master 節點啟動 Flink 叢集。
flink/bin/start-cluster.sh
初始化 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