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 :
- Récupérasion de la liste complète des workflows avec leurs codes
- Interrogation des planifications associées pour obtenir leur identifiant de调度
- 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.