Exécuter une tache PDI sur un serveur Carte depuis Apache Airflow

Nous avons changé d'ordonnanceur pour migrer vers Apache Airflow. Il nous fallait trouver comment exécuter nos taches pdi que nous comptons en dizaines. Plusieurs solutions sont apparues, un plugin pentaho data integration pour Airflow mais avec un suivi incertains, exécuter la tache dans un docker pdi mais aucune image officielle n'existe, exécuter la tache via ssh, installer pdi sur le serveur Airflow pour exécuter la tache localement ou pour terminer, exécuter la tache via le serveur carte et son api.
Nous verrons dans cet article comment configurer un serveur carte et comment via apache airflow exécuter une tache pentaho data integration sur le serveur carte sans aucun plugin.
Serveur Carte
L'exécution des taches pdi se fait sur un serveur windows dédié.
Configuration
la configuration du serveur carte pdi se fait via un fichier config_carte.xml.
<slave_config>
<slaveserver>
<name>ServerProd</name>
<hostname>HOSTNAME_SERVEUR</hostname>
<port>8282</port>
<username>user</username>
<password>password</password>
</slaveserver>
<max_log_lines>3000</max_log_lines>
<max_log_timeout_minutes>3000</max_log_timeout_minutes>
<object_timeout_minutes>3000</object_timeout_minutes>
</slave_config>
Démarrage
Le démarrage du serveur carte se fait via un simple script lancement_carte_94.bat. (sous windows)
E:\pdi-ce-9.4\data-integration\Carte.bat E:\pdi-ce-9.4\data-integration\config_carte.xml
Accéder à l'interface web du serveur carte.
http://HOSTNAME_SERVEUR:8282/kettle/status/

Api Serveur Carte
Nous allons utiliser plusieurs endpoint pour gérer l'execution d'une tache pdi.
- http://HOSTNAME_SERVEUR:8282/kettle/executeJob/?job=${PDI_JOB_PATH}&level=${PDI_LOG_LEVEL}
- http://HOSTNAME_SERVEUR:8282/kettle/jobStatus/?name=${PDI_JOB_NAME}&id=${PDI_JOB_ID}&xml=Y
Le endpoint "executeJob" permet de lancer l'exécution du job sur le serveur carte et de récupérer l'id du job.
Le endpoint jobStatus permet de récupérer les informations d'exécution. Notamment le status et le log final.
Apache Airflow
Nous utilisons un script pour lancer les taches pdi sur le serveur carte. Ce script est une adaptation d'un article sur le site : https://diethardsteiner.github.io/pdi/2020/04/01/Scheduling-a-PDI-Job-on-Apache-Airflow.html .
Notre version n'utilise pas de librairie xml car nous n'avions pas à disposition cette librairie mais le script utilise plutôt des regex en natif pour extraire les données nécessaires. Il y a également des ajouts pour gérer des cas particuliers et surtout le script est paramétrable pour passer des arguments depuis le dag.
Il reste a implémenter la possibilité de stopper la tache en cours si celle ci dépasse une durée paramétrée. Stopper le dag par un timeout coté Airflow ne servirait pas à grand chose puisque la tache va continuer à s'exécuter sur le serveur carte. Il faut donc envoyer l'ordre au serveur carte d'arrêter une tache avec un id spécifique.
MAJ 12/09/2025 : une nouvelle version (v2) est disponible avec la gestion d'un TIMEOUT en minute(s) pour stopper la tache pdi directement sur le serveur carte avant de générer une erreur pour Airflow. Les scripts en v2 sont disponibles.
MAJ 17/09/2025 : une nouvelle version (v3) est disponible avec la possibilité de donner des paramètres à la tache pdi. Les scripts en v3 sont disponibles.
Le script envoi un code d'erreur "exit 1" pour Apache Airflow si la tache pdi est en échec.
Il est tout à fait possible de créer des fonctions pythons pour faire la même chose.
Script
script : dags/scripts/pdi_carte_job_v3.sh
#!/bin/bash
# Passage de parametre depuis le dag
CARTE_USER=$1
CARTE_PASSWORD=$2
CARTE_HOSTNAME=$3
CARTE_PORT=$4
PDI_JOB_PATH=$5
PDI_JOB_NAME=$6
PDI_LOG_LEVEL=$7
SLEEP_INTERVAL_SECONDS=$8
LOG_TYPE=$9
TIMEOUT=${10}
PARAMS=${11}
# contrôle du Timeout
START_TIME=$(date +%s)
TIMEOUT_SECONDS=$((TIMEOUT * 60))
echo "Log type : "$LOG_TYPE
echo "Timeout paramétré : "$TIMEOUT" minutes"
echo "Timeout en secondes : "$TIMEOUT_SECONDS" secondes"
# parametres pour la tache PDI
if [[ -n "$PARAMS" ]]; then
echo "Paramètres fournis : "$PARAMS""
else
echo "Paramètres non fournis"
fi
#pattern pour les regex
PATTERN_JOB_ID="(.*?)<\/id>"
PATTERN_JOB_STATUS="(.*?)<\/status_desc>"
PATTERN_JOB_LOG="(.*?)<\/result>"
PATTERN_JOB_LOG_HTML="
Dag
Ci-dessous un exemple de dag pour lancer une tache pdi sur un serveur carte avec le script présenté ci-dessus.
dags/dag_pdi_carte_v3.py
from datetime import datetime, timedelta
import pendulum
from airflow import DAG
from airflow.operators.bash import BashOperator
import env.config as config
import env.env as env
default_args = {
'owner' : 'nicolas-dupont',
'description' : 'Test Carte PDI Curl',
'depend_on_past' : False,
'start_date' : pendulum.today(config.TIMEZONE).add(days=-1), # annee / mois / jour
'email_on_failure' : False,
'email_on_retry' : False,
'retries' : 0,
'retry_delay' : timedelta(minutes=5),
'dagrun_timeout' : timedelta(minutes=5),
}
with DAG('bash_curl_carte_pdi',
dag_display_name='Test Carte PDI Curl',
default_args=default_args,
catchup=False,
tags=["pdi","test","carte","9.4","remote"],
schedule_interval=None,
is_paused_upon_creation=True,
max_active_runs=1,
on_failure_callback=None,
) as dag:
BashPdiCarteJob = {'USERNAME':env.PDI_SERVEUR_USER,
'PASSWORD':env.PDI_SERVEUR_PASSWORD,
'HOSTNAME': env.PDI_SERVEUR,
'PORT':env.PDI_SERVEUR_PORT,
'JOB_PATH':f'{env.PDI_BASE_PATH}test/test_pdi_carte.kjb',
'JOB_NAME':'test_pdi_carte',
'LOG_LEVEL':'Basic',
'SLEEP_INTERVAL':'5',
'LOG_RETURN':'FULL' #full or pdi
'TIMEOUT':'1', # en minute
'PARAMS','param1=value1' #param1=value1¶m2=value2¶m3=value3
}
curl_carte_pdi = BashOperator(task_id='curl_carte_pdi',
bash_command=f'bash /dags/scripts/pdi_carte_job_v3.sh {BashPdiCarteJob['USERNAME']} {BashPdiCarteJob['PASSWORD']} {BashPdiCarteJob['HOSTNAME']} {BashPdiCarteJob['PORT']} {BashPdiCarteJob['JOB_PATH']} {BashPdiCarteJob['JOB_NAME']} {BashPdiCarteJob['LOG_LEVEL']} {BashPdiCarteJob['SLEEP_INTERVAL']} {BashPdiCarteJob['LOG_RETURN']} {BashPdiCarteJob['TIMEOUT']} {BashPdiCarteJob['PARAMS']} '
)
curl_carte_pdi
Tache Pdi
On exécute une tache très simple pour tester le fonctionnement. Une temporisation de 15 secondes, un calcul javascript pour créer des variables et pour terminer on log les variables.

Log
Sur Le Serveur Carte

Dans Apache Airflow

Vous noterez que le log final est obtenu seulement à la fin de l'exécution de la tache. Si on veut du détail durant l'exécution, il faut se connecter au serveur carte via http. C'est une limitation qui pourrait gêner pour une tache ayant un temps d'exécution important. Dans ces cas la, j'utilise plutôt un lancement via ssh qui permet d'avoir en live le log. Mais entre nous, sauf pour du débug, on ne s'amuse pas souvent à regarder ce qui se passe à chaque lancement d'un job.
Résumé
Nous utilisons ce script en production depuis quelques mois pour les taches pdi qui n'ont pas migré vers Apache Hop. Nous n'avons rencontré pour le moment aucun problème de fiabilité. Il faudra tout de même le faire évoluer pour pouvoir stopper une tache si un temps d'exécution défini est dépassé.
Retrouver les fichiers d'exemples sur mon github : https://github.com/NicoDupont/Resources/tree/master/ndl-blog/pdi_carte_airflow