Ce livre blanc explique comment utiliser l'outil de test de référence Nexmark pour évaluer les performances de traitement de flux de Realtime Compute for Apache Flink.
Aperçu des performances
Nexmark est une norme industrielle pour l'évaluation des performances des moteurs de traitement de flux. Il comprend 19 requêtes standard couvrant des scénarios typiques tels que le filtrage, l'agrégation, les jointures et les fenêtres. Cette rubrique utilise l'outil de test Nexmark pour évaluer complètement les performances de Realtime Compute for Apache Flink avec une configuration de 8 UC et une base de référence de 100 millions d'enregistrements en entrée pour chaque requête. Les résultats des tests montrent que :
Les requêtes simples, telles que q0, q1 et q2, atteignent un débit (RPS) de 4 à 6,5 millions d'enregistrements par seconde.
Les requêtes complexes d'agrégation et de fenêtrage, telles que q4, q5 et q16, affichent un débit compris entre 150 000 et 630 000 enregistrements par seconde.
Dans l'ensemble, Realtime Compute for Apache Flink offre des performances Nexmark 3,24 fois supérieures à celles de Flink open source.
Outil de test
Nexmark est une suite de tests de référence standard pour les moteurs de traitement de flux. Le modèle de test est le suivant :
Table source Nexmark : génère des données de test (événements Person, Auction et Bid) à un TPS spécifié.
Transformations : 19 requêtes Nexmark standard couvrant des scénarios typiques tels que le filtrage, la transformation, l'agrégation, les jointures et les fenêtres.
Table de destination blackhole : écrit les données dans une table de destination blackhole afin d'éliminer toute interférence liée au stockage externe, permettant ainsi d'évaluer uniquement la capacité de traitement du moteur Flink lui-même.
L'outil de test Nexmark utilisé dans cette rubrique s'appuie sur l'OpenAPI de Realtime Compute for Apache Flink. Il automatise l'intégralité du flux de travail, y compris la création des jobs, leur déploiement, leur surveillance et la collecte des résultats. Vous n'avez pas besoin d'écrire manuellement du code SQL ni de créer des jobs dans la console.
Environnement de test
Les jobs Flink de ce test ont utilisé les configurations d'optimisation suivantes :
|
Paramètre |
Valeur |
Description |
|
table.exec.mini-batch.enabled |
true |
Active l'agrégation Mini-Batch. |
|
table.exec.mini-batch.allow-latency |
2s |
Intervalle de mise en mémoire tampon Mini-Batch. |
|
table.optimizer.distinct-agg.split.enabled |
true |
Active l'optimisation de division pour l'agrégation Distinct. |
|
execution.checkpointing.interval |
3min |
Intervalle de checkpoint. |
Prérequis
Le kit de développement Java (JDK) version 1.8.x ou ultérieure est installé.
Vous avez activé Realtime Compute for Apache Flink et créé un espace de travail. Pour plus d'informations, consultez la rubrique Activation de Realtime Compute for Apache Flink.
Vous avez obtenu l'ID AccessKey et le secret AccessKey de votre compte Alibaba Cloud.
Procédure
Étape 1 : Télécharger l'outil de test
Téléchargez et extrayez le package de l'outil de test Nexmark nexmark-flink.tar.gz.
La structure des répertoires après extraction est la suivante :
nexmark-flink/
├── run_nexmark.sh # Test entry script
├── nexmark_env.sh # Environment variable configuration file (requires editing)
├── bin/ # Runtime scripts
├── conf/ # Flink job configurations
├── lib/ # JAR files (to be uploaded to the console)
└── queries-vvp/ # Nexmark Query SQL files
Étape 2 : Charger le fichier JAR Nexmark
Connectez-vous à la console Realtime Compute for Apache Flink.
Cliquez sur l'espace projet cible. Dans le volet de navigation de gauche, sélectionnez .
Sélectionnez et chargez le fichier
nexmark-flink-0.2-SNAPSHOT.jar. Ce fichier se trouve dans le répertoirenexmark-flink/libde l'outil de test.-
Une fois le chargement terminé, cliquez sur le nom du fichier pour copier son adresse OSS. Vous aurez besoin de cette adresse lors d'une étape de configuration ultérieure. Le format du chemin d'accès au fichier varie selon le type de stockage :
-
Stockage OSS Bucket :
oss://<OSS Bucket name>/artifacts/namespaces/<project space name>/<file name>Exemple :
oss://oss-test/artifacts/namespaces/flink-default/nexmark-flink-0.2-SNAPSHOT.jar -
Stockage entièrement géré :
oss://flink-fullymanaged-<workspace ID>/artifacts/namespaces/<project space name>/<file name>Exemple : oss://flink-fullymanaged-e6a123456789/artifacts/namespaces/flink-default/nexmark-flink-0.2-SNAPSHOT.jar
Pour afficher le type de stockage de votre espace de travail, accédez à la console de gestion Realtime Compute for Apache Flink, localisez l'espace de travail cible et cliquez sur Actions dans la colonne Details.
-
Étape 3 : Configurer les paramètres d'exécution
Modifiez le fichier nexmark-flink/nexmark_env.sh et définissez les paramètres suivants.
|
Paramètre |
Description |
Exemple |
|
END_POINT |
L'endpoint pour Realtime Compute for Apache Flink. Sélectionnez l'endpoint correspondant à votre région. Pour plus d'informations, consultez la rubrique Endpoints. |
ververica.cn-hangzhou.aliyuncs.com |
|
AK |
L'ID AccessKey de votre compte Alibaba Cloud. |
- |
|
SK |
Le secret AccessKey de votre compte Alibaba Cloud. |
- |
|
WORK_SPACE |
L'ID de votre espace de travail. |
e6a123456789 |
|
NAMESPACE |
Le nom de votre espace projet. |
flink-default |
|
NEXMARK_JAR |
L'adresse OSS du fichier JAR que vous avez chargé à l'étape 2. |
oss://flink-fullymanaged-e6a123456789/artifacts/namespaces/flink-default/nexmark-flink-0.2-SNAPSHOT.jar |
|
FLINK_VERSION |
La version du moteur Flink à tester. |
vvr-11.5-jdk11-flink-1.20 |
|
QUERIES |
Spécifiez les requêtes à exécuter. Séparez plusieurs requêtes par des virgules, par exemple |
all |
L'exécution de toutes les requêtes prend beaucoup de temps. Chaque requête doit passer par des étapes telles que la création du job, la génération des données et l'exécution des calculs. Nous vous recommandons de commencer par exécuter une seule requête (par exemple, en définissant QUERIES sur q0) afin de vérifier que la configuration de l'environnement et les paramètres sont corrects avant de lancer un test à grande échelle.
Étape 4 : Exécuter le test
-
Dans le répertoire
nexmark-flink, exécutez la commande suivante../run_nexmark.sh L'outil de test crée et exécute automatiquement les jobs Nexmark via l'OpenAPI.
-
Une fois le test terminé, la durée de chaque requête s'affiche en millisecondes. L'exemple ci-dessous montre un résultat type :
INFO com.github.nexmark.flink.vvp.Nexmark - q0 13078 ============================================================================ ✓ Benchmark execution completed successfully ============================================================================
Résultats de performance
Le tableau suivant compare les performances Nexmark entre Flink open source (1.20.4) et Realtime Compute for Apache Flink (vvr-11.5-jdk11-flink-1.20) sur une configuration de 8 UC. Chaque requête traite 100 millions d'enregistrements en entrée. RPS = Nombre d'enregistrements en entrée ÷ Durée.
Les données de test suivantes ont été collectées dans un environnement matériel spécifique et avec des versions de moteur particulières. Les performances réelles peuvent varier en raison des mises à niveau matérielles et des mises à jour du moteur. Ces résultats sont fournis à titre indicatif uniquement.
|
Requête |
Flink open source sur ECS Version : 1.20.4 |
Realtime Compute for Apache Flink Version : vvr-11.5-jdk11-flink-1.20 |
|||
|
Durée (ms) |
RPS |
Durée (ms) |
RPS |
RPS par rapport à l'open source (×) |
|
|
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 |
|
Total |
2383209 |
49,695,131 |
9550361 |
15,317,480 |
3.24 |
Procédure de test pour Flink open source
Pour reproduire les résultats sur Flink open source déployé sur ECS, suivez cette procédure.
Préparation de l'environnement
Créez un cluster Flink à l'aide d'EMR sur ECS avec la configuration suivante :
Version EMR : EMR-5.21.0
Spécifications matérielles : trois instances ecs.g6a.xlarge (4 vCPU / 16 Go), comprenant un nœud maître et deux nœuds core.
Activez les services Hadoop et HDFS.
-
Configurez la connexion sans mot de passe entre tous les nœuds. Par exemple, téléchargez votre fichier de clé privée (tel que
key.pem) sur le nœud maître. Ajoutez ensuite la configuration suivante au fichier~/.ssh/configsur le nœud maître. Remplacez les adresses IP et le chemin d'accès au fichier par vos valeurs réelles.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 noUtilisez ssh pour vérifier que la connexion sans mot de passe entre les nœuds fonctionne correctement. Si une erreur « bad permissions » se produit, exécutez
chmod 600 /path/to/key.pempour corriger les autorisations.
Préparation logicielle
-
Téléchargez le package Flink cible (Apache Flink Downloads) et le package de test Nexmark (nexmark-flink.tgz). Chargez les packages sur le nœud maître et extrayez-les.
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
-
Copiez les fichiers JAR du répertoire
nexmark/libversflink/lib. Ces JAR contiennent le générateur de données Nexmark.cp nexmark/lib/* flink/lib/
-
Définissez les variables d'environnement. Modifiez
~/.bashrc, ajoutez la configuration suivante, puis exécutezsource ~/.bashrcpour appliquer les modifications.Configurez les chemins d'accès en fonction de votre environnement réel.
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)
Configuration et démarrage du cluster
-
Configurez les workers Flink. Ce test utilise huit TaskManagers, déployés comme suit : deux sur le nœud maître et trois sur chaque nœud core.
Modifiez
flink/conf/workers, en veillant à remplacer l'adresse IP par la valeur réelle.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
-
Remplacez
flink/conf/config.yamlparnexmark/conf/config.yaml, et mettez à jour les éléments de configuration suivants :jobmanager.rpc.address: L'adresse IP du nœud maître, par exemple192.168.0.0.state.checkpoints.dir: Le chemin HDFS, par exemplehdfs:///checkpointstaskmanager.memory.process.size:4G
Modifiez
nexmark/conf/nexmark.yamlet définisseznexmark.metric.reporter.hostsur l'adresse IP du nœud maître.-
Distribuez les répertoires
flinketnexmarkainsi que les configurations des variables d'environnement vers chaque nœud core.Remplacez les adresses IP par vos valeurs réelles.
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:~/Une fois la distribution terminée, exécutez
source ~/.bashrcsur chaque nœud core pour activer les variables d'environnement.
-
Sur le nœud maître, démarrez le cluster Flink.
flink/bin/start-cluster.sh
-
Initialisez l'environnement de test Nexmark. Ce script configure le Metric Reporter requis sur chaque nœud.
nexmark/bin/setup_cluster.sh
Définir les limites de ressources
Après le démarrage du cluster Flink, vous devez utiliser cgroups pour limiter l'utilisation du CPU de chaque processus TaskManager à 75 %. Cela empêche les TaskManagers de perdre la connexion en raison de délais d'expiration des signaux de maintien de liaison (heartbeat) causés par la contention des ressources.
Exécutez les commandes suivantes sur tous les nœuds où des TaskManagers sont en cours d'exécution, y compris le nœud maître :
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
Ici, cpu.cfs_quota_us / cpu.cfs_period_us = 300000 / 100000 = 3, ce qui signifie que le cgroup peut utiliser un maximum de 3 cœurs CPU (75 % de 4 vCPU).
Dans l'environnement entièrement géré de Realtime Compute for Apache Flink, vous n'avez pas besoin de définir manuellement une limite d'utilisation des ressources. Le service exploite pleinement les ressources de calcul que vous avez achetées.
Exécuter Nexmark
Sur le nœud maître, exécutez la commande suivante et attendez qu'elle se termine pour afficher les résultats :
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