Skip to main content
Antes de que existiera el provider oficial de ClickHouse, la mayoría de los usuarios de Airflow se conectaban a ClickHouse mediante el package comunitario airflow-clickhouse-plugin. Esta guía explica cómo migrar un despliegue existente a 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 de common.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 y apache-airflow-providers-common-sql 1.32.0 o posterior. Si utiliza una release anterior, actualice primero Airflow.
Los dos paquetes residen en espacios de nombres de Python distintos, por lo que pueden instalarse en paralelo mientras migra DAG por DAG. Elimine el plugin cuando ya nada lo importe:

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 de extra 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:
Después:
Verifique cada connection migrada antes de modificar los DAG:

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:
Después:

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_pull posterior para que tome el último elemento.
  • Pasa las sentencias como una única cadena separada por ; y establece split_statements=True. Con el valor predeterminado return_last=True, el operator enviará entonces solo las filas de la última sentencia, igual que el plugin.
Un ClickHouseOperator que insertaba una lista de filas mediante parameters pasa a ser un SQLInsertRowsOperator:
Pasa siempre 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:
Después:
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:
Después:
Si tu invocable necesita la fila completa, pasa 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.sql y elimina el prefijo ClickHouse del nombre de la class.
  • Pasa conn_id de forma explícita. Las envolturas interpretaban un conn_id ausente o None como clickhouse_default; las clases de common.sql no tienen valor predeterminado y fallan si no se indica. default_args={"conn_id": "clickhouse_default"} cubre todo un DAG.
  • database= en los operadores y hook_params={"schema": ...} en el sensor siguen funcionando; el hook del provider trata schema como un alias de database.
  • ClickHouseDbApiHook pasa a ser ClickHouseHook. Su argumento de constructor schema se sigue aceptando como alias; database es 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 devolviera clickhouse-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

  1. Airflow es 2.11 o posterior.
  2. El puerto HTTP es accesible desde los workers; los certificados TLS son válidos para el endpoint HTTP.
  3. Cada ClickHouse connection: tipo clickhouse, puerto 8123 o 8443, extras traducidos según el paso 2 y verificados con airflow connections test.
  4. Imports reemplazados según el paso 3; prefijos ClickHouse eliminados de las envolturas de common.sql.
  5. clickhouse_conn_id renombrado a conn_id en los operators y sensors, y conn_id establecido en cada task que dependía del valor predeterminado del plugin.
  6. settings= trasladado a hook_params={"session_settings": ...}.
  7. hook.execute("INSERT ... VALUES", rows) reemplazado por bulk_insert_rows; ClickHouseOperator(parameters=rows) reemplazado por SQLInsertRowsOperator.
  8. Callables de los sensors ajustados de result[0][0] al valor de la celda sin más.
  9. Usos de with_column_types, external_tables, columnar, query_id y types_check reescritos con un handler o get_client().
  10. Revisados los consumidores posteriores de XComs procedentes de tasks con múltiples sentencias y de INSERT.
  11. Actualizado el código que captura excepciones de clickhouse_driver.
  12. airflow-clickhouse-plugin y clickhouse-driver desinstalados.

Uso de un asistente de programación con IA

La correspondencia anterior es deliberadamente mecánica, de modo que un asistente de programación pueda aplicarla a un repositorio de DAG. Un prompt que ha dado buenos resultados:
Revise el diff. Los cambios en la connection, la semántica de los sensores y los consumers de XCom son los puntos donde las reescrituras automatizadas suelen fallar.
Última modificación el 26 de septiembre de 2026