apache-airflow-providers-clickhousedb. Para conocer el funcionamiento del
provider en sí, consulta Conectar Apache Airflow a ClickHouse.
¿Por qué un provider nativo?
Los providers son la forma en que Airflow se integra con sistemas de terceros. Se publican, se prueban y se documentan junto con el resto del ecosistema de Airflow, y de este en concreto se encarga la comunidad de Airflow junto con el ClickHouse team. Está construido sobre clickhouse-connect, el Python client que ClickHouse desarrolla y mantiene, y no sobre un driver mantenido por la comunidad. Así, las nuevas funcionalidades y correcciones del servidor llegan a los usuarios de Airflow por una vía con soporte. Pasar al provider te aporta un package con un hogar oficial, los operators y sensores estándar decommon.sql y un tipo de connection que aparece
en la UI de Airflow igual que el de cualquier otra base de datos.
Los dos packages se diferencian en algo más que en las rutas de import. El plugin se comunica con ClickHouse mediante el
protocolo TCP nativo usando clickhouse-driver. El provider lo hace por HTTP(S) usando
clickhouse-connect y se integra con los operators genéricos de common.sql en lugar de incluir
operators específicos de ClickHouse.
Lee la guía completa una vez antes de cambiar nada. El cambio de connection, en particular, afecta
a todos los DAG a la vez.
De un vistazo
Paso 1: comprobar los prerequisites e instalar
El provider requiere Airflow 2.11 o posterior yapache-airflow-providers-common-sql 1.32.0 o
posterior. Si utiliza una release anterior, actualice primero Airflow.
Paso 2: Actualizar las conexiones
Este es el paso que hace que todo falle si lo omite. Todas las conexiones de ClickHouse existentes apuntan al puerto nativo, y el provider necesita el puerto HTTP.
Si ClickHouse se encuentra detrás de un firewall o de un balanceador de carga, asegúrese de que el puerto HTTP sea accesible desde
los workers antes de realizar el cambio. Compruebe que la interfaz HTTP esté habilitada en el servidor (
http_port o
https_port en la configuración del servidor). ClickHouse Cloud expone HTTPS únicamente en
8443.
Las conexiones almacenadas como URI (clickhouse://user:pass@host:9000/db?secure=true) ya tienen el
tipo de conexión clickhouse, porque Airflow lo deduce del scheme del URI. En su caso solo cambian el puerto y las
claves adicionales; los valores de la cadena de consulta se analizan como JSON, por lo que secure=true sigue siendo un
valor booleano.
Extras de conexión
El plugin pasaba cada clave deextra directamente a clickhouse_driver.Client, por lo que las conexiones pueden
llevar cualquier keyword argument de clickhouse-driver. El provider solo lee un conjunto fijo de claves y
reenvía todo lo demás mediante client_kwargs. Haz la equivalencia de la siguiente manera:
Antes:
Paso 3: Reemplazar los imports
Las envolturas de
common.sql con el prefijo ClickHouse solo existían para inyectar el hook del plugin. El
provider registra el tipo de connection clickhouse, de modo que las clases de common.sql sin prefijo resuelven
el hook a partir de la connection por sí solas. Si usabas esas envolturas, la migración se reduce normalmente a cambiar la
línea de import, quitar el prefijo ClickHouse y pasar conn_id de forma explícita (véase el
paso 7).
Paso 4: de ClickHouseOperator a SQLExecuteQueryOperator
Antes:
Resultados de múltiples sentencias
El plugin enviaba a XCom el resultado de la última sentencia.SQLExecuteQueryOperator devuelve
un resultado por sentencia cuando sql es una lista, por lo que el ejemplo anterior envía [[], [(12345.0,)]]
mientras que el plugin enviaba [(12345.0,)]. Elige una de estas opciones:
- Modifica el
xcom_pullposterior para que tome el último elemento. - Pasa las sentencias como una única cadena separada por
;y establecesplit_statements=True. Con el valor predeterminadoreturn_last=True, el operator enviará entonces solo las filas de la última sentencia, igual que el plugin.
ClickHouseOperator que insertaba una lista de filas mediante parameters pasa a ser un
SQLInsertRowsOperator:
columns; sin él, el operator busca la table mediante SQLAlchemy.
Conservar los tipos de columna
with_column_types=True devolvía (rows, [(name, type), ...]). Puede reproducirlo con un handler; el cursor de
clickhouse-connect expone los nombres de tipo de ClickHouse en cursor.description:
Paso 5: de ClickHouseHook.execute a los métodos de DbApiHook
El hook del plugin exponía un único método, execute, que replicaba clickhouse_driver.Client.execute. El
hook del provider es un DbApiHook, por lo que dispone de los métodos estándar que tiene cualquier otro provider SQL:
run, get_records, get_first, get_pandas_df, get_df, insert_rows y test_connection.
Los argumentos del constructor clickhouse_conn_id y database no cambian.
fetch_all_handler y los demás handler se importan desde
airflow.providers.common.sql.hooks.handlers.
El uso más habitual del hook en el código de la época del plugin es el bulk insert. No puede hacerse con una simple llamada a run,
porque el cursor DB-API intentaría formatear las filas dentro de la cadena SQL. Utiliza en su lugar el insert nativo:
Antes:
bulk_insert_rows requiere column_names. batch_size es opcional y limita el uso de memoria con
entradas muy grandes. El método genérico insert_rows(table, rows, target_fields=[...], executemany=True)
también termina en un insert nativo, pero solo con executemany=True; por defecto envía un HTTP request por
fila.
Para todo lo que la superficie DB-API no cubre, get_client() devuelve el client sin procesar de clickhouse-connect,
configurado a partir de la connection de Airflow. Es el reemplazo de todos los argumentos
específicos de clickhouse-driver que exponía el plugin:
Paso 6: de ClickHouseSensor a SqlSensor
Este es el único reemplazo en el que cambia la entrada del invocable.
Como el plugin entregaba el result set completo, el código de sensor existente lo indexa. Elimine
esa indexación al migrar:
Antes:
selector=lambda row: row. Si necesita todas las filas, escribe
la comprobación en SQL para que la consulta devuelva un único valor booleano o un recuento.
Paso 7: La familia de envolturas common.sql
El código que usaba ClickHouseSQLExecuteQueryOperator, ClickHouseSqlSensor y las demás
envolturas con prefijo ClickHouse es el que requiere menos trabajo:
- Cambia el import al módulo
common.sqly elimina el prefijoClickHousedel nombre de la class. - Pasa
conn_idde forma explícita. Las envolturas interpretaban unconn_idausente oNonecomoclickhouse_default; las clases decommon.sqlno tienen valor predeterminado y fallan si no se indica.default_args={"conn_id": "clickhouse_default"}cubre todo un DAG. database=en los operadores yhook_params={"schema": ...}en el sensor siguen funcionando; el hook del provider trataschemacomo un alias dedatabase.ClickHouseDbApiHookpasa a serClickHouseHook. Su argumento de constructorschemase sigue aceptando como alias;databasees la forma nativa de ClickHouse y tiene precedencia cuando se indican ambos.- La conexión sigue necesitando los cambios de puerto y de extras del paso 2. Las envolturas también usaban el protocolo nativo.
Diferencias de comportamiento a revisar
Incluso después de que el código compile, hay algunos aspectos que se comportan de forma distinta en tiempo de ejecución. Valor de XCom de las tareas INSERT. El plugin enviaba lo que devolvieraclickhouse-driver, que en el caso de una inserción VALUES con parámetros era el número de filas insertadas. El provider envía un conjunto de resultados vacío para las sentencias que no devuelven filas. Las tareas posteriores que lean el número de filas desde XCom deberán obtenerlo de otra manera, por ejemplo con un SELECT count() posterior.
Settings de sesión frente a sentencias SET. Ambos paquetes ejecutan una lista de varias sentencias sobre una única conexión, y clickhouse-connect crea una sesión por cliente de forma predeterminada, por lo que una sentencia SET al principio de la lista debería seguir aplicándose a las sentencias posteriores. Aun así, conviene usar session_settings: es explícito y admite plantillas, y funciona igual tanto si el servidor o un proxy intermedio mantiene la sesión como si no. Confirma el comportamiento en tu entorno si tus DAG dependen de SET.
Excepciones. Los errores ahora son clickhouse_connect.driver.exceptions.DatabaseError, OperationalError o ProgrammingError en lugar de clickhouse_driver.errors.ServerException y NetworkError. Actualiza las cláusulas except, el código de on_failure_callback y la lógica de reintentos que inspeccione los tipos de excepción.
Compresión. clickhouse-driver dejaba la compresión desactivada salvo que se estableciera compression; clickhouse-connect habilita la compresión de las respuestas HTTP de forma predeterminada y negocia el algoritmo con el servidor. Establece "compress": false en el campo extra de la conexión para restaurar el comportamiento anterior. El paquete clickhouse-cityhash que el plugin necesitaba para la compresión nativa ya no hace falta; lz4 sigue instalado como dependencia de clickhouse-connect.
Correspondencia de tipos. Ambos drivers devuelven tipos nativos de Python, pero se trata de bases de código distintas. Revisa las tareas que dependan de tipos exactos para DateTime64 con zonas horarias, Decimal, UUID, columnas Nullable y valores anidados de Array o Map, sobre todo cuando el resultado se envía a XCom y se consume más adelante.
Identificación de consultas en system.query_log. Las consultas ahora llegan a través de la interfaz HTTP, por lo que aparecen con interface = 2 en lugar de 1, y la columna http_user_agent de system.query_log incluye las versiones de Airflow y del provider, además del extra client_name si está definido. Habrá que actualizar cualquier monitorización que filtrara por el protocolo nativo o por el nombre de cliente de clickhouse-driver. Un SELECT que no devuelve filas produce una segunda entrada: el cursor DB-API ejecuta SELECT * FROM (...) LIMIT 0 para recuperar los metadatos de las columnas.
Timeouts sobre HTTP. send_receive_timeout es ahora el timeout de lectura HTTP, y cualquier proxy o balanceador de carga situado entre los workers y ClickHouse aplica su propio idle timeout a la petición. Las sentencias que se ejecutaban durante muchos minutos sobre el protocolo nativo pueden requerir que se eleven esos límites.
Gestión de conexiones. El hook crea un cliente de clickhouse-connect por cada llamada a run o get_client, igual que el plugin abría una nueva conexión nativa por cada execute. Los clientes comparten un grupo de conexiones HTTP a nivel de proceso, por lo que llamar a close() en un cliente obtenido con get_client() es una buena práctica más que un requisito; el grupo se libera cuando el proceso de la tarea termina. El cliente es un gestor de contexto, así que with hook.get_client() as client: es la forma más limpia.
Lista de comprobación
- Airflow es 2.11 o posterior.
- El puerto HTTP es accesible desde los workers; los certificados TLS son válidos para el endpoint HTTP.
- Cada ClickHouse connection: tipo
clickhouse, puerto8123o8443, extras traducidos según el paso 2 y verificados conairflow connections test. - Imports reemplazados según el paso 3; prefijos
ClickHouseeliminados de las envolturas decommon.sql. clickhouse_conn_idrenombrado aconn_iden los operators y sensors, yconn_idestablecido en cada task que dependía del valor predeterminado del plugin.settings=trasladado ahook_params={"session_settings": ...}.hook.execute("INSERT ... VALUES", rows)reemplazado porbulk_insert_rows;ClickHouseOperator(parameters=rows)reemplazado porSQLInsertRowsOperator.- Callables de los sensors ajustados de
result[0][0]al valor de la celda sin más. - Usos de
with_column_types,external_tables,columnar,query_idytypes_checkreescritos con unhandleroget_client(). - Revisados los consumidores posteriores de XComs procedentes de tasks con múltiples sentencias y de INSERT.
- Actualizado el código que captura excepciones de
clickhouse_driver. airflow-clickhouse-pluginyclickhouse-driverdesinstalados.