Skip to main content
Avant l’arrivée du provider ClickHouse officiel, la plupart des utilisateurs d’Airflow se connectaient à ClickHouse via le paquet communautaire airflow-clickhouse-plugin. Ce guide détaille la migration d’un déploiement existant vers apache-airflow-providers-clickhousedb. Pour comprendre le fonctionnement du provider lui-même, consultez Connecter Apache Airflow à ClickHouse.

Pourquoi un provider natif ?

Les providers sont le moyen par lequel Airflow s’intègre à des systèmes tiers. Ils sont publiés, testés et documentés en même temps que le reste de l’écosystème Airflow, et celui-ci est maintenu par la communauté Airflow en collaboration avec l’équipe ClickHouse. Il repose sur clickhouse-connect, le client Python que ClickHouse développe et prend en charge lui-même, plutôt que sur un driver maintenu par la communauté. Les nouvelles fonctionnalités du serveur et les correctifs parviennent donc aux utilisateurs d’Airflow par une voie prise en charge. Passer au provider vous offre un paquet disposant d’un point d’ancrage officiel, des operators et sensors common.sql standard, ainsi que d’un type de connexion qui apparaît dans l’UI d’Airflow comme pour toute autre base de données. Les deux paquets diffèrent par bien plus que les chemins d’import. Le plugin communique avec ClickHouse via le protocole TCP natif au moyen de clickhouse-driver. Le provider communique en HTTP(S) avec clickhouse-connect et s’appuie sur les operators common.sql génériques au lieu de fournir des operators spécifiques à ClickHouse.
Lisez l’intégralité du guide avant d’effectuer la moindre modification. Le changement de connexion, en particulier, affecte tous les DAG en même temps.

En bref

Étape 1 : vérifier les prerequisites et installer

Le provider nécessite Airflow 2.11 ou une version plus récente, ainsi que apache-airflow-providers-common-sql 1.32.0 ou une version plus récente. Commencez par effectuer l’upgrade d’Airflow si vous utilisez une release plus ancienne.
Les deux paquets résident dans des namespaces Python différents : ils peuvent donc être installés côte à côte pendant que vous migrez vos DAG un par un. Supprimez le plugin dès que plus rien ne l’importe :

Étape 2 : mettre à jour les connexions

C’est l’étape qui casse tout si vous l’oubliez. Toutes les connexions ClickHouse existantes pointent vers le port natif, alors que le provider a besoin du port HTTP. Si ClickHouse se trouve derrière un firewall ou un répartiteur de charge, assurez-vous que le port HTTP est joignable depuis les workers avant de basculer. Vérifiez que l’interface HTTP est activée sur le serveur (http_port ou https_port dans la configuration du serveur). ClickHouse Cloud n’expose HTTPS que sur le port 8443. Les connexions stockées sous forme d’URI (clickhouse://user:pass@host:9000/db?secure=true) possèdent déjà le type de connexion clickhouse, car Airflow le déduit du scheme de l’URI. Pour celles-ci, seuls le port et les clés supplémentaires changent ; les valeurs de la chaîne de requête sont interprétées comme du JSON, si bien que secure=true reste un booléen.

Extras de connexion

Le plugin transmettait chaque clé de extra directement à clickhouse_driver.Client ; les connexions peuvent donc contenir n’importe quel keyword argument de clickhouse-driver. Le provider ne lit qu’un ensemble fixe de clés et transmet tout le reste via client_kwargs. Procédez à la conversion comme suit : Avant :
Après :
Vérifiez chaque connexion migrée avant de toucher aux DAG :

Étape 3 : remplacer les imports

Les wrappers common.sql préfixés par ClickHouse n’existaient que pour injecter le hook du plugin. Le provider enregistre le type de connexion clickhouse : les classes common.sql sans préfixe résolvent donc elles-mêmes le hook à partir de la connexion. Si vous utilisiez ces wrappers, la migration se limite en général à la ligne d’import, à la suppression du préfixe ClickHouse et au passage explicite de conn_id (voir l’étape 7).

Étape 4 : de ClickHouseOperator à SQLExecuteQueryOperator

Avant :
Après :

Résultats multi-statements

Le plugin poussait le résultat du dernier statement dans XCom. SQLExecuteQueryOperator renvoie un résultat par statement lorsque sql est une liste ; l’exemple ci-dessus pousse donc [[], [(12345.0,)]] là où le plugin poussait [(12345.0,)]. Choisissez l’une de ces options :
  • Modifier le xcom_pull en aval pour prendre le dernier élément.
  • Passer les statements sous forme d’une chaîne unique séparée par des ; et définir split_statements=True. Avec la valeur par défaut return_last=True, l’operator ne pousse alors que les lignes du dernier statement, ce qui correspond au comportement du plugin.
Un ClickHouseOperator qui insérait une liste de lignes via parameters devient un SQLInsertRowsOperator :
Passez toujours columns ; sinon, l’operator recherche la table via SQLAlchemy.

Conserver les types de colonnes

with_column_types=True renvoyait (rows, [(name, type), ...]). Reproduisez ce comportement à l’aide d’un handler ; le curseur clickhouse-connect expose les noms de types ClickHouse dans cursor.description :

Étape 5 : de ClickHouseHook.execute aux méthodes de DbApiHook

Le hook du plugin n’exposait qu’une seule méthode, execute, calquée sur clickhouse_driver.Client.execute. Le hook du provider est un DbApiHook : il dispose donc des méthodes standard communes à tous les autres providers SQL : run, get_records, get_first, get_pandas_df, get_df, insert_rows et test_connection. Les arguments de constructeur clickhouse_conn_id et database restent inchangés. fetch_all_handler et les autres handlers s’importent depuis airflow.providers.common.sql.hooks.handlers. Dans le code de l’ère du plugin, l’usage le plus courant du hook est le bulk insert. Il ne peut pas se réduire à un simple appel à run, car le curseur DB-API tenterait de formater les lignes dans la chaîne SQL. Utilisez plutôt le native insert : Avant :
Après :
bulk_insert_rows requiert column_names. batch_size est optionnel et limite la mémoire utilisée sur de très gros volumes en entrée. La méthode générique insert_rows(table, rows, target_fields=[...], executemany=True) aboutit également à un insert natif, mais uniquement avec executemany=True ; par défaut, une requête HTTP est envoyée par ligne. Pour tout ce que l’interface DB-API ne couvre pas, get_client() renvoie le client clickhouse-connect brut configuré à partir de la connexion Airflow. Il remplace l’ensemble des arguments spécifiques à clickhouse-driver qu’exposait le plugin :

Étape 6 : ClickHouseSensor vers SqlSensor

C’est le seul remplacement où l’entrée du callable change. Comme le plugin transmettait l’intégralité du result set, le code de sensor existant y accède par index. Supprimez cette indexation lors de la migration : Avant :
Après :
Si votre callable a besoin de la ligne entière, passez selector=lambda row: row. S’il a besoin de toutes les lignes, écrivez le check en SQL de sorte que la requête renvoie un seul booléen ou un décompte.

Étape 7 : la famille de wrappers common.sql

Le code qui utilisait ClickHouseSQLExecuteQueryOperator, ClickHouseSqlSensor et les autres wrappers préfixés par ClickHouse demande le moins de travail :
  • Modifiez l’import pour pointer vers le module common.sql et supprimez le préfixe ClickHouse du nom de la classe.
  • Passez explicitement conn_id. Les wrappers considéraient un conn_id absent ou à None comme clickhouse_default ; les classes common.sql n’ont aucune valeur par défaut et échouent s’il n’est pas fourni. default_args={"conn_id": "clickhouse_default"} couvre l’ensemble d’un DAG.
  • database= sur les opérateurs et hook_params={"schema": ...} sur le sensor continuent de fonctionner ; le hook du provider traite schema comme un alias de database.
  • ClickHouseDbApiHook devient ClickHouseHook. Son argument de constructeur schema est toujours accepté comme alias ; database est l’orthographe native de ClickHouse et a la préséance lorsque les deux sont fournis.
  • La connexion nécessite toujours les modifications de port et d’extras de l’étape 2 : les wrappers utilisaient eux aussi le protocole natif.

Différences de comportement à examiner

Même après une compilation réussie du code, certains éléments se comportent différemment à l’exécution. Valeur XCom des tasks INSERT. Le plugin transmettait ce que renvoyait clickhouse-driver, soit, pour un insert VALUES avec paramètres, le nombre de lignes insérées. Le provider transmet un ensemble de résultats vide pour les statements qui ne renvoient aucune ligne. Les tasks en aval qui lisent le nombre de lignes depuis XCom doivent l’obtenir autrement, par exemple avec un SELECT count() complémentaire. Paramètres de session contre statements SET. Les deux paquets exécutent une liste de plusieurs statements sur une seule connexion, et clickhouse-connect crée par défaut une session par client ; un statement SET placé en début de liste devrait donc toujours s’appliquer aux statements suivants. Privilégiez malgré tout session_settings : c’est explicite et templatisé, et cela fonctionne que le serveur ou un proxy intermédiaire conserve ou non la session. Vérifiez le comportement dans votre environnement si vos DAG dépendent de SET. Exceptions. Les erreurs sont désormais clickhouse_connect.driver.exceptions.DatabaseError, OperationalError ou ProgrammingError au lieu de clickhouse_driver.errors.ServerException et NetworkError. Mettez à jour les clauses except, le code on_failure_callback et la logique de reprise qui inspecte les types d’exception. Compression. clickhouse-driver laissait la compression désactivée sauf si compression était défini ; clickhouse-connect active par défaut la compression des réponses HTTP et négocie l’algorithme avec le serveur. Définissez "compress": false dans l’extra de connexion pour rétablir l’ancien comportement. Le paquet clickhouse-cityhash dont le plugin avait besoin pour la compression native n’est plus requis ; lz4 reste installé en tant que dépendance de clickhouse-connect. Correspondance des types. Les deux drivers renvoient des types Python natifs, mais il s’agit de bases de code différentes. Examinez les tasks qui dépendent de types exacts pour DateTime64 avec fuseaux horaires, Decimal, UUID, les colonnes Nullable et les valeurs Array ou Map imbriquées, en particulier lorsque le résultat est transmis à XCom puis consommé en aval. Identification des requêtes dans system.query_log. Les requêtes arrivent désormais via l’HTTP interface, si bien qu’elles apparaissent avec interface = 2 au lieu de 1, et la colonne http_user_agent de system.query_log contient les versions d’Airflow et du provider ainsi que l’extra client_name s’il est défini. Toute supervision filtrant sur le native protocol ou sur le nom de client clickhouse-driver doit être mise à jour. Un SELECT qui ne renvoie aucune ligne produit une seconde entrée : le curseur DB-API exécute SELECT * FROM (...) LIMIT 0 pour récupérer les métadonnées des colonnes. Timeouts en HTTP. send_receive_timeout correspond maintenant au timeout de lecture HTTP, et tout proxy ou répartiteur de charge situé entre les workers et ClickHouse applique son propre idle timeout à la requête. Les statements qui s’exécutaient pendant de longues minutes via le native protocol peuvent nécessiter un relèvement de ces limites. Gestion des connexions. Le hook crée un client clickhouse-connect par appel à run ou get_client, à l’image du plugin qui ouvrait une nouvelle native connection à chaque execute. Les clients partagent un HTTP connection pool à l’échelle du processus ; appeler close() sur un client issu de get_client() relève donc de la bonne hygiène plutôt que de l’obligation : le pool est libéré à la fin du processus de la task. Le client est un context manager : with hook.get_client() as client: est donc la forme la plus propre.

Checklist

  1. Airflow est en version 2.11 ou plus récente.
  2. Le port HTTP est joignable depuis les workers ; les certificats TLS sont valides pour l’endpoint HTTP.
  3. Chaque connexion à ClickHouse : type clickhouse, port 8123 ou 8443, extras convertis selon l’étape 2, vérifiés avec airflow connections test.
  4. Les imports sont remplacés selon l’étape 3 ; les préfixes ClickHouse sont supprimés des wrappers common.sql.
  5. clickhouse_conn_id est renommé en conn_id sur les operators et les sensors, et conn_id est défini sur chaque task qui s’appuyait sur la valeur par défaut du plugin.
  6. settings= est déplacé dans hook_params={"session_settings": ...}.
  7. hook.execute("INSERT ... VALUES", rows) est remplacé par bulk_insert_rows ; ClickHouseOperator(parameters=rows) est remplacé par SQLInsertRowsOperator.
  8. Les callables des sensors passent de result[0][0] à la valeur brute de la cell.
  9. Les utilisations de with_column_types, external_tables, columnar, query_id et types_check sont réécrites avec un handler ou get_client().
  10. Les consumers en aval des XComs issus des tasks multi-statement et INSERT sont passés en revue.
  11. Le code qui intercepte les exceptions clickhouse_driver est mis à jour.
  12. airflow-clickhouse-plugin et clickhouse-driver sont désinstallés.

Utilisation d’un assistant de codage IA

Le mapping ci-dessus est volontairement mécanique, de sorte qu’un assistant de codage puisse l’appliquer à un repository de DAG. Voici un prompt qui a donné de bons résultats :
Examinez le diff. Les modifications de connexion, la sémantique des sensors et les consommateurs XCom sont les points sur lesquels les réécritures automatisées échouent.
Dernière modification le 26 septembre 2026