PARTAGE

Le 16 avril 2025 | Update 18 septembre 2025 | 8 mins read

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.

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=""
CARTE_SERVER_URL="http://${CARTE_USER}:${CARTE_PASSWORD}@${CARTE_HOSTNAME}:${CARTE_PORT}"

# construction de l'url pour lancer la tache pdi



URL="${CARTE_SERVER_URL}/kettle/executeJob/?job=file:///${PDI_JOB_PATH}&level=${PDI_LOG_LEVEL}"
# Ajouter le paramètre optionnel s'il est fourni
if [[ -n "$PARAMS" ]]; then
    URL="${URL}&${PARAMS}"
fi
echo "Url envoyée au server pdi : "
echo $URL

# commande curl vers le serveur pdi
PDI_JOB_RES=$(curl -s "${URL}" | tr -d "\r\n")
#echo "Reponse du server pdi : " $PDI_JOB_RES

if [[ $PDI_JOB_RES =~ $PATTERN_JOB_ID ]]; then
    PDI_JOB_ID=${BASH_REMATCH[1]}
else
    PDI_JOB_ID="no found job id"
fi
#echo "recherche id : " $BASH_REMATCH
#echo "The PDI job ID is: " $PDI_JOB_ID

function getPDIJobStatus {
    STATUS=$(curl -s "${CARTE_SERVER_URL}/kettle/jobStatus/?name=${PDI_JOB_NAME}&id=${PDI_JOB_ID}&xml=Y" | tr -d "\r\n")
    #echo 'resultat status job : '$STATUS
    if [[ $STATUS =~ $PATTERN_JOB_STATUS ]]; then
        JOB_STATUS=${BASH_REMATCH[1]}
    else
        JOB_STATUS="no found job status"
    fi
    echo $JOB_STATUS
}

#log du server
function getPDIJobFullLog {
    LOG=$(curl -s "${CARTE_SERVER_URL}/kettle/jobStatus/?name=${PDI_JOB_NAME}&id=${PDI_JOB_ID}&xml=Y")
    echo $LOG
}

#Stop de la tache
function stopPdiJob {
    curl -s "${CARTE_SERVER_URL}/kettle/stopJob/?name=${PDI_JOB_NAME}&id=${PDI_JOB_ID}&xml=Y"
    echo "--------------------------------------"
    echo "Tache PDI arrété -> timeout dépassé !!"
}

#log de la tache + transformation
function getPDIJobFullLogHtml {
    LOG=$(curl -s "${CARTE_SERVER_URL}/kettle/jobStatus/?name=${PDI_JOB_NAME}&id=${PDI_JOB_ID}")
    if [[ $LOG =~ $PATTERN_JOB_LOG_HTML ]]; then
        HTML_LOG=${BASH_REMATCH[0]}
    else
        HTML_LOG="no found job status"
    fi
    echo $HTML_LOG
}

#log de la transformation ???
function getPDIJobLog {
    LOG=$(curl -s "${CARTE_SERVER_URL}/kettle/jobStatus/?name=${PDI_JOB_NAME}&id=${PDI_JOB_ID}&xml=Y" | tr -d "\r\n")
    echo 'resultat log job : '$LOG
    if [[ $LOG =~ $PATTERN_JOB_LOG ]]; then
        JOB_LOG=${BASH_REMATCH[1]}
    else
        JOB_LOG="no found job log"
    fi
    echo $JOB_LOG
}

echo "Tache PDI : " $PDI_JOB_PATH
echo "PDI job ID : " $PDI_JOB_ID
echo "PDI server url : ${CARTE_SERVER_URL}/kettle/status/" 
echo "PDI server url pour ce job : ${CARTE_SERVER_URL}/kettle/jobStatus/?name=${PDI_JOB_NAME}&id=${PDI_JOB_ID}" 

# loop as long as the job is running
PDI_JOB_STATUS=$(getPDIJobStatus)
echo "Status du Job -> " $PDI_JOB_STATUS
while [[ $PDI_JOB_STATUS == "Running" || $PDI_JOB_STATUS == "Waiting" ]]
do
    # Récupération du status de la tache => pour ne pas killer en timeout si terminée
    PDI_JOB_STATUS=$(getPDIJobStatus)

    # Vérification durée d'exécution
    CURRENT_TIME=$(date +%s)
    ELAPSED_TIME=$((CURRENT_TIME - START_TIME))
    if [[ $ELAPSED_TIME -gt $TIMEOUT_SECONDS && ${PDI_JOB_STATUS} != 'Finished' ]]; then
        stopPdiJob
        PDI_JOB_STATUS="TIMEOUT" # on sort de la boucle avec cette valeur au prochain test
    fi

    echo "The PDI job status is: ${PDI_JOB_STATUS} -> Durée exécution : $((ELAPSED_TIME / 60)) minutes et $((ELAPSED_TIME % 60)) secondes"
    #echo "Temps écoulé : $((ELAPSED_TIME / 60)) minutes et $((ELAPSED_TIME % 60)) secondes"
    sleep ${SLEEP_INTERVAL_SECONDS}
done 

# afficher log xml ou html complet ou seulement la partie pdi
if [[ $LOG_TYPE == 'FULL' ]]; then
    echo "Full log ..."
    echo ""
    #echo $(getPDIJobFullLog)
    echo $(getPDIJobFullLogHtml)
else
    echo "Pdi log ..."
    echo ""
    echo $(getPDIJobLog)
fi

sleep 1

# Envoyer une erreur à airflow si status nok ou timeout
if [[ ${PDI_JOB_STATUS} == 'Finished' ]]; then
    echo "Task PDI Finished -> exit 0"
    exit 0
elif [[ ${PDI_JOB_STATUS} == 'TIMEOUT' ]]; then
    echo "!! Task PDI en Timeout -> exit 1 !!"
    exit 1
else
    echo "!! Task PDI en Erreur -> exit 1 !!"
    exit 1
fi

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