Automatisation de l'import et de la mise en ligne de workflows DolphinScheduler via API

Documantation des API

DolphinScheduler expose une documentation Swagger accessible à l'adresse suivante :

http://VOTRE_SERVEUR:12345/dolphinscheduler/swagger-ui/index.html?language=zh_CN&lang=cn

La documentation reste succincte et nécessite une exploration personnelle des différents points de terminaison.

Authentification par jeton

Toutes les requêtes API requièrent un jeton d'authentification. Créez-en un via Centre de sécurité → Gestion des jetons dans l'interface d'administration. Ce jeton doit être inclus dans l'en-tête HTTP de chaque appel.

api_token = 'VOTRE_JETON'
http_headers = {
    'Accept': 'application/json',
    'token': api_token
}

L'identifiant du projet (project_id) est visible dans l'URL du navigateur lors de la consultation des workflows d'un projet.

Import de workflows

Le point de terminaison dédié à l'import est :

endpoint_import = 'http://VOTRE_SERVEUR:12345/dolphinscheduler/projects/{project_id}/process-definition/import'

L'import exige un transfert binaire du fichier. Un workflow exporté manuellement depuis l'interface peut servir de modèle pour construire les fichiers à importer.

import requests

def push_workflow_definition(chemin_fichier, url_import, headers):
    with open(chemin_fichier, 'rb') as fh:
        payload = {'file': fh}
        resultat = requests.post(url_import, headers=headers, files=payload)
        if resultat.status_code != 200:
            print(f"Échec d'import : {chemin_fichier} (HTTP {resultat.status_code})")
            return False
        return True

L'appel répété de cette fonction sur une collection de fichiers permet l'import massif.

Mise en ligne des workflows

Après l'import, chaque workflow doit être activé manuellement par défaut. L'API permet d'automatiser cette opération selon la séquence suivante :

  1. Récupérasion de la liste complète des workflows avec leurs codes
  2. Interrogation des planifications associées pour obtenir leur identifiant de调度
  3. Activation de chaque planification hors ligne

Récupération paginée des workflows

import json

def recuperer_tous_workflows(url_liste, headers):
    collecte = []
    page_courante = 1
    taille_page = 10
    
    while True:
        requete = f'{url_liste}?pageSize={taille_page}&pageNo={page_courante}&searchVal='
        reponse = requests.get(requete, headers=headers)
        
        if reponse.status_code != 200:
            print(f'Erreur HTTP {reponse.status_code} : {reponse.text}')
            break
        
        contenu = json.loads(reponse.content.decode())
        total_global = contenu['data']['total']
        page_donnees = contenu['data']['totalList']
        
        collecte.extend(page_donnees)
        
        if page_courante * taille_page >= total_global:
            break
        
        page_courante += 1
    
    return collecte

Extraction des codes uniques de chaque workflow :

liste_complete = recuperer_tous_workflows(jobs_url, http_headers)
codes_workflow = [entree['code'] for entree in liste_complete]

Obtention des identifiants de planification

endpoint_schedules = 'http://VOTRE_SERVEUR:12345/dolphinscheduler/projects/{project_id}/schedules?pageSize=1000&pageNo=1&processDefinitionCode='

def trouver_id_planification(code_wf, url_schedules, headers):
    requete_complete = url_schedules + str(code_wf)
    reponse = requests.get(requete_complete, headers=headers)
    
    if reponse.status_code != 200:
        return None
    
    donnees = json.loads(reponse.content.decode())
    items_planif = donnees['data']['totalList']
    
    if items_planif and items_planif[0]['releaseState'] == 'OFFLINE':
        return items_planif[0]['id']
    
    return None

Ce filtre ignore les planifications déjà actives.

Activation en ligne

endpoint_online = 'http://VOTRE_SERVEUR:12345/dolphinscheduler/projects/{project_id}/schedules/{schedule_id}/online'

def activer_workflow(id_planif, url_online, headers):
    url_final = url_online.replace('{schedule_id}', str(id_planif))
    reponse = requests.post(url_final, headers=headers)
    
    if reponse.status_code == 200:
        print(f'Workflow {id_planif} activé avec succès')
        return True
    else:
        print(f'Échec activation pour ID {id_planif}')
        return False

En enchaînant ces fonctions, l'intégralité du cycle import → activation s'exécute sans intervention manuelle.

Étiquettes: DolphinScheduler API REST automatisation workflow orchestration

Publié le 24 août à 14h26