Questions fréquemment posées sur l'utilisation des clusters DataFlow.
-
Utilisation et opérations du cluster :
Comment soumettre un job à un cluster DataFlow depuis une machine externe ?
Comment résoudre les noms d'hôte d'un cluster DataFlow depuis une machine externe ?
Comment accéder au Flink HistoryServer dans un cluster DataFlow ?
Comment utiliser les connecteurs commerciaux pris en charge par un cluster DataFlow ?
Comment activer la haute disponibilité pour le JobManager d'un job Flink ?
-
Problèmes liés aux jobs :
Comment résoudre les problèmes liés au stockage amont et aval (connecteurs) ?
Où se trouvent les journaux du client et comment les consulter ?
Où se trouvent les journaux du cluster et comment les consulter ?
Que faire si un package JAR de job entre en conflit avec le package JAR Flink du cluster ?
Comment activer les graphiques en flammes (flame graphs) pour un job Flink ?
Comment résoudre l'erreur « java.lang.OutOfMemoryError: GC overhead limit exceeded » ?
Consulter les journaux du cluster
Consultez les journaux en fonction du statut du JobManager :
Si le JobManager s'est arrêté, récupérez les journaux sur votre machine locale en exécutant la commande yarn logs -applicationId application_xxxx_yy. Vous pouvez également consulter les journaux en accédant au lien des journaux du job terminé dans l'interface web YARN.
-
Si le JobManager est toujours en cours d'exécution, utilisez l'une des méthodes suivantes :
Accédez à l'interface web Flink pour consulter les journaux.
Utilisez des outils en ligne de commande. Exécutez yarn logs -applicationId application_xxxx_yy -am ALL -logFiles jobmanager.log pour afficher les journaux du JobManager, ou exécutez yarn logs -applicationId application_xxxx_yy -containerId container_xxxx_yy_aa_bb -logFiles taskmanager.log pour afficher les journaux du TaskManager.
Résoudre les conflits de packages JAR
Ce problème provoque généralement des erreurs telles que NoSuchFieldError/NoSuchMethodError/ClassNotFoundException dans les journaux du job. Pour diagnostiquer et résoudre ce problème, suivez les étapes ci-dessous :
Identifiez la classe de dépendance en conflit. En vous basant sur la classe d'exception indiquée dans le message d'erreur, repérez le package JAR de dépendance qui contient cette classe. Ensuite, dans le répertoire où se trouve le fichier pom.xml de votre job, exécutez mvn dependency:tree pour afficher l'arborescence des dépendances et déterminer leur origine.
-
Excluez la classe de dépendance en conflit.
Si le champ d'application (scope) du package JAR est incorrectement défini dans le fichier pom.xml, modifiez le scope en
providedpour exclure le package JAR.Si vous devez impérativement utiliser le package JAR contenant la classe d'exception, vous pouvez ajouter une règle d'exclusion pour supprimer la classe conflictuelle spécifique.
Si vous devez utiliser la classe d'exception et ne pouvez pas la remplacer par la version correspondante du cluster, utilisez le plugin Maven Shade pour isoler (« shadier ») la classe.
De plus, si plusieurs versions d'un package JAR existent dans le classpath, la version de la classe utilisée par le job dépend de l'ordre de chargement des classes. Pour confirmer à partir de quel package JAR une classe spécifique est chargée, vous pouvez définir le paramètre JVM env.java.opts: -verbose:class dans le fichier flink-conf.yaml ou spécifier le paramètre dynamique -Denv.java.opts="-verbose:class" afin d'imprimer les classes chargées et leurs sources.
RemarquePour un JobManager ou un TaskManager, ces informations sont imprimées dans le fichier
jobmanager.outoutaskmanager.out.
Soumettre des jobs depuis une machine externe
Pour soumettre un job à un cluster DataFlow depuis une machine externe :
Assurez-vous que la machine externe peut se connecter au cluster DataFlow via le réseau.
-
Configurez l'environnement Hadoop YARN sur la machine cliente qui soumet le job Flink.
Dans un cluster DataFlow, le logiciel Hadoop YARN est installé dans le répertoire /opt/apps/YARN/yarn-current et ses fichiers de configuration se trouvent dans le répertoire /etc/taihao-apps/hadoop-conf/. Vous devez télécharger le répertoire yarn-current et le répertoire hadoop-conf sur la machine cliente.
Ensuite, configurez les variables d'environnement suivantes sur la machine cliente.
export HADOOP_HOME=/path/to/yarn-current && \ export PATH=${HADOOP_HOME}/bin/:$PATH && \ export HADOOP_CLASSPATH=$(hadoop classpath) && \ export HADOOP_CONF_DIR=/path/to/hadoop-confImportantLes fichiers de configuration Hadoop, tels que yarn-site.xml, utilisent un nom de domaine complet (FQDN) pour les adresses de service comme le ResourceManager. Par exemple,
master-1-1.c-xxxxxxxxxx.cn-hangzhou.emr.aliyuncs.com. Si vous soumettez des jobs depuis une machine externe, assurez-vous que ces FQDN peuvent être résolus, ou remplacez les FQDN par leurs adresses IP correspondantes dans les fichiers de configuration. Une fois la configuration terminée, démarrez un job Flink sur la machine externe. Par exemple, exécutez la commande
flink run -d -t yarn-per-job -ynm flink-test $FLINK_HOME/examples/streaming/TopSpeedWindowing.jar. Vous pouvez alors voir le job Flink correspondant dans l'interface web YARN du cluster DataFlow.
Résoudre les noms d'hôte du cluster depuis une machine externe
Utilisez l'une des méthodes suivantes pour résoudre les noms d'hôte du cluster DataFlow depuis une machine externe :
Modifiez le fichier /etc/hosts sur la machine cliente pour ajouter les mappages entre les noms d'hôte et les adresses IP.
-
Utilisez le service DNS fourni par Alibaba Cloud DNS PrivateZone.
Si vous disposez de votre propre service de résolution de noms de domaine, vous pouvez également configurer les paramètres d'exécution JVM suivants pour l'utiliser.
env.java.opts.client: "-Dsun.net.spi.nameservice.nameservers=xxx -Dsun.net.spi.nameservice.provider.1=dns,sun -Dsun.net.spi.nameservice.domain=yyy"
Vérifier le statut d'un job Flink
-
Utilisez la console EMR.
EMR prend en charge Knox, ce qui vous permet d'accéder aux interfaces web des services tels que YARN et Flink via Internet. Vous pouvez accéder à l'interface web Flink via YARN. Pour plus d'informations, consultez Afficher le statut du job sur l'interface web.
Utilisez un tunnel SSH. Pour plus d'informations, consultez Créer un tunnel SSH pour accéder aux interfaces web des composants open source.
-
Accédez directement à l'API REST YARN.
curl --compressed -v -H "Accept: application/json" -X GET "http://master-1-1:8088/ws/v1/cluster/apps?states=RUNNING&queue=default&user.name=***"RemarqueAssurez-vous que votre groupe de sécurité autorise l'accès aux ports 8443 et 8088 pour atteindre l'API REST YARN. Alternativement, assurez-vous que le cluster DataFlow et votre nœud client se trouvent dans le même Virtual Private Cloud (VPC).
Accéder aux journaux d'un job Flink
Pour un job en cours d'exécution, vous pouvez accéder à ses journaux via l'interface web Flink.
Pour un job terminé, vous pouvez consulter ses statistiques sur le Flink HistoryServer ou accéder à ses journaux en exécutant la commande
yarn logs -applicationId application_xxxx_yyyy. Les journaux des jobs terminés sont stockés par défaut dans le répertoire hdfs:///tmp/logs/$USERNAME/logs/ du cluster HDFS.
Accéder au Flink HistoryServer
Par défaut, un cluster DataFlow exécute un Flink HistoryServer sur le nœud master-1-1 (la première machine du groupe de serveurs maîtres) sur le port 18082. Le HistoryServer collecte les statistiques des jobs terminés. Pour y accéder :
Configurez une règle de groupe de sécurité pour autoriser l'accès au port 18082 sur le nœud
master-1-1.Accédez directement à http://$master-1-1-ip:18082.
Le Flink HistoryServer ne stocke pas les journaux détaillés des jobs terminés. Pour consulter les journaux, utilisez l'API YARN ou l'interface web YARN.
Utiliser les connecteurs commerciaux
Les clusters DataFlow incluent des connecteurs commerciaux pour Hologres, SLS, MaxCompute, DataHub, Elasticsearch, ClickHouse, etc. Vous pouvez les utiliser conjointement avec les connecteurs open source. L'exemple suivant montre comment utiliser le connecteur Hologres.
-
Développement du job
-
Téléchargez le package JAR du connecteur commercial depuis le cluster DataFlow (situé dans le répertoire /opt/apps/FLINK/flink-current/opt/connectors). Ensuite, installez le connecteur dans votre environnement Maven local en exécutant la commande suivante.
mvn install:install-file -Dfile=/path/to/ververica-connector-hologres-1.13-vvr-4.0.7.jar -DgroupId=com.alibaba.ververica -DartifactId=ververica-connector-hologres -Dversion=1.13-vvr-4.0.7 -Dpackaging=jar -
Ajoutez la dépendance suivante au fichier pom.xml de votre projet.
<dependency> <groupId>com.alibaba.ververica</groupId> <artifactId>ververica-connector-hologres</artifactId> <version>1.13-vvr-4.0.7</version> <scope>provided</scope> </dependency>
-
-
Exécution du job
-
Méthode 1 :
-
Copiez le connecteur Hologres dans un répertoire distinct.
hdfs mkdir hdfs:///flink-current/opt/connectors/hologres/ hdfs cp hdfs:///flink-current/opt/connectors/ververica-connector-hologres-1.13-vvr-4.0.7.jar hdfs:///flink-current/opt/connectors/hologres/ververica-connector-hologres-1.13-vvr-4.0.7.jar -
Lors de la soumission du job, ajoutez le paramètre suivant à la commande.
-D yarn.provided.lib.dirs=hdfs:///flink-current/opt/connectors/hologres/
-
-
Méthode 2 :
Copiez le connecteur Hologres dans le répertoire /opt/apps/FLINK/flink-current/opt/connectors/ververica-connector-hologres-1.13-vvr-4.0.7.jar sur le client de soumission de job. Cette structure de répertoire doit correspondre à celle du cluster DataFlow.
-
Lors de la soumission du job, ajoutez le paramètre suivant à la commande.
-C file:///opt/apps/FLINK/flink-current/opt/connectors/ververica-connector-hologres-1.13-vvr-4.0.7.jar
Méthode 3 : Intégrez le connecteur Hologres au package JAR de votre job.
-
Utiliser GeminiStateBackend
Les clusters DataFlow utilisent par défaut le GeminiStateBackend de niveau entreprise, qui offre des performances 3 à 5 fois supérieures à celles des backends d'état open source. Pour plus d'informations sur les configurations avancées de GeminiStateBackend, consultez Configurations du backend d'état de niveau entreprise.
Utiliser un backend d'état open source
Les clusters DataFlow utilisent par défaut GeminiStateBackend. Pour utiliser un backend d'état open source tel que RocksDB pour un job spécifique, spécifiez-le avec l'indicateur -D :
flink run-application -t yarn-application -D state.backend=rocksdb /opt/apps/FLINK/flink-current/examples/streaming/TopSpeedWindowing.jar
Alternativement, pour appliquer cette modification à tous les jobs suivants, accédez à la console EMR, changez la valeur du paramètre state.backend vers le backend d'état souhaité (par exemple, rocksdb). Cliquez sur Save, puis cliquez sur Deploy Client Configuration.
Consulter les journaux du client
La variable d'environnement FLINK_LOG_DIR spécifie le répertoire des journaux du client Flink. Sa valeur par défaut est /var/log/taihao-apps/flink (avant la version 3.43.0, la valeur par défaut était /mnt/disk1/log/flink). Pour consulter les journaux complets du client, tels que les journaux SQL Client, vérifiez ce répertoire.
Paramètres de job non pris en compte
Lorsque vous exécutez un job Flink depuis la ligne de commande, placez les paramètres du job après le package JAR du job Flink. Par exemple : flink run -d -t yarn-per-job test.jar arg1 arg2.
Résoudre l'erreur « Multiple factories... »
-
Cause
Le classpath contient plusieurs implémentations d'un connecteur. Cela se produit généralement lorsqu'une dépendance de connecteur est empaquetée dans le JAR du job alors que le même connecteur existe également dans le répertoire $FLINK_HOME/lib, provoquant un conflit de dépendances.
-
Solution
La solution consiste à supprimer la dépendance en double. Pour des étapes de dépannage détaillées, consultez Que faire si un package JAR de job entre en conflit avec le package JAR Flink du cluster ?
Activer la haute disponibilité du JobManager
Les clusters DataFlow exécutent les jobs Flink en mode YARN. Pour activer la haute disponibilité (HA) pour le JobManager, suivez le guide de Configuration de la communauté. Voici un exemple de configuration.
high-availability: zookeeper
high-availability.zookeeper.quorum: 192.168.**.**:2181,192.168.**.**:2181,192.168.**.**:2181
high-availability.zookeeper.path.root: /flink
high-availability.storageDir: hdfs:///flink/recovery
Après avoir activé la haute disponibilité, le JobManager redémarre au maximum une fois en cas de défaillance par défaut. Si vous souhaitez que le JobManager redémarre plusieurs fois, vous devez également définir le paramètre yarn.resourcemanager.am.max-attempts de YARN et le paramètre yarn.application-attempts de Flink. Pour plus d'informations, consultez la documentation officielle Apache Flink. Selon l'expérience, vous devriez également augmenter la valeur du paramètre yarn.application-attempt-failures-validity-interval de la valeur par défaut de 10 000 millisecondes (10 secondes) à une valeur plus élevée, telle que 300 000 millisecondes (5 minutes), afin d'éviter que le JobManager ne redémarre continuellement.
Consulter les métriques d'un job Flink
Dans la console EMR, accédez à la page Monitoring du cluster cible et cliquez sur Metric Monitoring.
Dans la liste déroulante Dashboard, sélectionnez FLINK.
-
Sélectionnez l'ID d'application et l'ID de job pour le job que vous souhaitez consulter. Les métriques de surveillance du job s'affichent alors.
RemarqueLes options ID d'application et ID de job sont disponibles uniquement si des jobs Flink sont en cours d'exécution dans le cluster.
Certaines métriques, telles que
sourceIdleTime, sont générées uniquement si la source et le puits (sink) correspondants sont configurés.
Résoudre les problèmes de connecteurs
Pour les questions courantes concernant le stockage amont et aval, consultez Connecteurs.
Résoudre les erreurs d'accès à OSS sans mot de passe
Gérez le problème en fonction du message d'erreur spécifique :
-
Message d'erreur :
java.lang.UnsupportedOperationException: Recoverable writers on Hadoop are only supported for HDFS.Cause : Les clusters DataFlow utilisent le JindoSDK intégré pour l'accès sans mot de passe à OSS et des API telles que StreamingFileSink. Aucune configuration supplémentaire issue de la documentation communautaire n'est requise. Son ajout peut provoquer un conflit de dépendances déclenchant cette erreur.
Solution : Sur la machine de soumission de job dans votre cluster, vérifiez la présence d'un répertoire oss-fs-hadoop dans le répertoire $FLINK_HOME/plugins. S'il existe, supprimez le répertoire et resoumettez le job.
-
Message d'erreur :
Could not find a file system implementation for scheme 'oss'. The scheme is directly supported by Flink through the following plugin: flink-oss-fs-hadoop. .....Cause : Dans EMR 3.40 et versions antérieures, les machines du groupe de serveurs maîtres autres que
master-1-1peuvent ne pas disposer des packages JAR liés à Jindo.-
Solution :
-
Pour EMR 3.40 et versions antérieures : Vérifiez si les packages JAR liés à Jindo, tels que jindo-flink-4.0.0-full.jar, existent dans le répertoire $FLINK_HOME/lib sur la machine de soumission de job. S'ils sont manquants, exécutez la commande suivante pour copier les packages JAR requis dans le répertoire $FLINK_HOME/lib, puis resoumettez le job.
cp /opt/apps/extra-jars/flink/jindo-flink-*-full.jar $FLINK_HOME/lib -
Pour les versions EMR postérieures à 3.40 :
Pour le mode Flink on YARN : Les versions plus récentes disposent d'un mécanisme optimisé pour la prise en charge d'OSS. Les jobs qui lisent et écrivent dans OSS peuvent s'exécuter normalement même si les packages JAR liés à Jindo ne sont pas présents dans le répertoire $FLINK_HOME/lib.
-
Pour les autres modes de déploiement : Vérifiez si les packages JAR liés à Jindo, tels que jindo-flink-4.0.0-full.jar, existent dans le répertoire $FLINK_HOME/lib de la machine de soumission de job. S'ils sont manquants, exécutez la commande suivante pour les copier dans le répertoire $FLINK_HOME/lib, puis resoumettez le job.
cp /opt/apps/extra-jars/flink/jindo-flink-*-full.jar $FLINK_HOME/lib
-
Résoudre l'erreur « TaskManager heartbeat timed out »
-
Cause
Un délai d'attente du heartbeat d'un TaskManager peut résulter de plusieurs causes. Vérifiez les journaux du TaskManager pour obtenir des messages d'erreur spécifiques. Les causes courantes incluent une insuffisance de mémoire heap ou une erreur de mémoire insuffisante (OOM) due à une fuite de mémoire dans le code du job. Pour plus d'informations, consultez Comment résoudre l'erreur « java.lang.OutOfMemoryError: GC overhead limit exceeded » ?.
-
Solution
Si vous rencontrez cette erreur, augmentez l'allocation de mémoire ou analysez l'utilisation de la mémoire du job pour diagnostiquer davantage le problème.
Résoudre l'erreur « GC overhead limit exceeded »
-
Cause
Le garbage collector (GC) consomme un temps excessif car la mémoire est insuffisante. Les causes courantes incluent une fuite de mémoire dans le code (comme dans une UDF) ou une allocation de mémoire trop faible par rapport aux besoins du job.
-
Solution
Avant de relancer le job, spécifiez le paramètre JVM suivant à l'aide de l'indicateur -D pour enregistrer un vidage de tas (heap dump) lorsqu'une OutOfMemoryError se produit :
-D env.java.opts="-XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/tmp/dump.hprof".Ajoutez le paramètre
env.java.opts: -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/tmp/dump.hprofau fichierflink-conf.yamlpour configurer les vidages de tas en cas d'OutOfMemoryError.
Une fois l'erreur reproduite, vous pouvez analyser le fichier de vidage de tas spécifié par
HeapDumpPathà l'aide d'outils tels que MAT ou jvisualvm pour déterminer la cause racine.
Zéro « Records Received » pour les jobs à opérateur unique
Ceci est normal. La métrique Records Received de Flink décrit la communication de données entre différents opérateurs. Lorsqu'un job est optimisé en un seul opérateur, cette métrique sera toujours égale à 0.
Activer les graphiques en flammes (flame graphs) pour les jobs Flink
Les graphiques en flammes visualisent la consommation de CPU à travers les méthodes au sein d'un processus, ce qui vous aide à identifier les goulets d'étranglement de performance. Flink prend en charge les graphiques en flammes depuis la version 1.13, mais la fonctionnalité est désactivée par défaut pour éviter d'affecter les jobs de production. Si vous devez utiliser les graphiques en flammes pour analyser les performances d'un job, accédez à l'onglet Configure du service Flink dans la console EMR. Dans le fichier flink-conf.yaml, ajoutez un nouvel élément de configuration avec le paramètre rest.flamegraph.enabled et définissez sa valeur sur true. Pour obtenir des instructions sur l'ajout d'un élément de configuration, consultez Gérer les éléments de configuration.
Pour plus d'informations sur les graphiques en flammes, consultez Flame Graphs.
Résoudre l'erreur « NoSuchFieldError: DEPLOYMENT_MODE »
-
Cause
Le package JAR de votre job inclut directement ou indirectement une dépendance
flink-coreincompatible avec la version Flink du cluster, provoquant un conflit de dépendances. -
Solution
Ajoutez la configuration suivante à votre fichier pom.xml pour définir le
scopede la dépendanceflink-coresurprovided. Cela résout le problème.<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-core</artifactId> <!-- Change to your own flink version --> <version>1.16.1</version> <scope>provided</scope> </dependency>RemarqueVous devez modifier la
versionpour correspondre à votre version Flink.Pour localiser plus précisément la source de cette dépendance, consultez Que faire si un package JAR de job entre en conflit avec le package JAR Flink du cluster ?.