Comment passer des arguments Java aux artefacts de travail Flink en mode Application
Je viens de mettre à jour le Flink de la version 1.10 à 1.11. Dans la version 1.11, Flink fournit de nouvelles fonctionnalités permettant aux utilisateurs de déployer le travail en mode Application sur Kubernetes.https://ci.apache.org/projects/flink/flink-docs-release-1.11/ops/deployment/kubernetes.html#deploy-session-cluster
Dans la V1.10, nous démarrons le cluster Flink K8s, puis soumettons le travail à Flink par exécution
exec ./bin/flink run \
-d \
/streakerflink_deploy.jar \
--arg1 blablabla
--arg2 blablabla
--arg3 blablabla
...
Nous passons les arguments java via cette commande.
Mais, dans la version 1.11, si nous exécutons le mode Application, nous n'avons pas besoin d'exécuter la flink runcommande ci-dessus. Je me demande comment passer les arguments au travail Flink en mode Application (alias Job Cluster)?
Toute aide serait appréciée!
Réponses
Puisque vous utilisez un graphique de barre pour démarrer le cluster Flink sur Kubernetes (alias K8s), je suppose que vous parlez du mode K8s autonome. En fait, le mode Application est très similaire au cluster de travaux de la version 1.10 et des versions antérieures. Vous pouvez donc définir les arguments du travail dans le argschamp jobmanager-job.yamlcomme suit.
...
args: ["standalone-job", "--job-classname", "org.apache.flink.streaming.examples.join.WindowJoin", "--windowSize", "3000", "--rate", "100"]
...
Si vous voulez vraiment parler du mode K8s natif, alors il pourrait être directement ajouté après la flink run-applicationcommande.
$ ./bin/flink run-application -p 8 -t kubernetes-application \
-Dkubernetes.cluster-id=<ClusterId> \
-Dtaskmanager.memory.process.size=4096m \
-Dkubernetes.taskmanager.cpu=2 \
-Dtaskmanager.numberOfTaskSlots=4 \
-Dkubernetes.container.image=<CustomImageName> \
local:///opt/flink/examples/streaming/WindowJoin.jar \
--windowSize 3000 --rate 100
Remarque: veuillez garder à l'esprit que la principale différence entre le mode K8 autonome et le mode K8 natif est l'allocation dynamique des ressources. En mode natif, nous avons un client K8s intégré, donc Flink JobManager peut allouer / libérer des pods TaskManager sur demande. Actuellement, le mode natif ne peut être utilisé que dans les commandes Flink ( kubernetes-session.sh, flink run-application).
Le mode d'application de Flink sur Kubernetes est décrit dans la documentation . Vous devez créer une image Docker contenant votre travail. Le travail peut être exécuté en utilisant ./bin/flink run-application [...]comme décrit dans la documentation.