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 sensorscommon.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 queapache-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.
É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é deextra 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 :
É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 :
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_pullen aval pour prendre le dernier élément. - Passer les statements sous forme d’une chaîne unique séparée par des
;et définirsplit_statements=True. Avec la valeur par défautreturn_last=True, l’operator ne pousse alors que les lignes du dernier statement, ce qui correspond au comportement du plugin.
ClickHouseOperator qui insérait une liste de lignes via parameters devient un
SQLInsertRowsOperator :
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 :
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 :
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.sqlet supprimez le préfixeClickHousedu nom de la classe. - Passez explicitement
conn_id. Les wrappers considéraient unconn_idabsent ou àNonecommeclickhouse_default; les classescommon.sqln’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 ethook_params={"schema": ...}sur le sensor continuent de fonctionner ; le hook du provider traiteschemacomme un alias dedatabase.ClickHouseDbApiHookdevientClickHouseHook. Son argument de constructeurschemaest toujours accepté comme alias ;databaseest 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 renvoyaitclickhouse-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
- Airflow est en version 2.11 ou plus récente.
- Le port HTTP est joignable depuis les workers ; les certificats TLS sont valides pour l’endpoint HTTP.
- Chaque connexion à ClickHouse : type
clickhouse, port8123ou8443, extras convertis selon l’étape 2, vérifiés avecairflow connections test. - Les imports sont remplacés selon l’étape 3 ; les préfixes
ClickHousesont supprimés des wrapperscommon.sql. clickhouse_conn_idest renommé enconn_idsur les operators et les sensors, etconn_idest défini sur chaque task qui s’appuyait sur la valeur par défaut du plugin.settings=est déplacé danshook_params={"session_settings": ...}.hook.execute("INSERT ... VALUES", rows)est remplacé parbulk_insert_rows;ClickHouseOperator(parameters=rows)est remplacé parSQLInsertRowsOperator.- Les callables des sensors passent de
result[0][0]à la valeur brute de la cell. - Les utilisations de
with_column_types,external_tables,columnar,query_idettypes_checksont réécrites avec unhandlerouget_client(). - Les consumers en aval des XComs issus des tasks multi-statement et INSERT sont passés en revue.
- Le code qui intercepte les exceptions
clickhouse_driverest mis à jour. airflow-clickhouse-pluginetclickhouse-driversont désinstallés.