Comment passer des arguments Java aux artefacts de travail Flink en mode Application

Sep 02 2020

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

1 YangWang Sep 04 2020 at 10:23

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).

1 mapohl Sep 02 2020 at 20:38

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.