You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

WSL环境Docker中Airflow无法连接ClickHouse查询问题

问题:Airflow连接ClickHouse时错误使用9000端口导致连接拒绝

问题场景

在WSL环境中通过Docker部署ClickHouse服务器与Airflow,两者均处于运行状态且已配置连接,但执行查询时Airflow日志出现连接拒绝错误:

[2024-06-04, 12:17:49 UTC] {logging_mixin.py:188} INFO - >>> Running 'etl' function. Logged at 2024-06-04 12:17:49.678034+00:00
.....
[2024-06-04, 12:17:53 UTC] {logging_mixin.py:188} INFO - >>> Running 'clickhouse_connect' function. Logged at 2024-06-04 12:17:53.590231+00:00
[2024-06-04, 12:17:53 UTC] {logging_mixin.py:188} INFO - Connecting to ClickHouse Cluster
[2024-06-04, 12:17:53 UTC] {base.py:84} INFO - Using connection ID 'my_clickhouse' for task execution.
[2024-06-04, 12:17:53 UTC] {logging_mixin.py:188} INFO - 8123
[2024-06-04, 12:17:53 UTC] {logging_mixin.py:188} INFO - Connected successfully
[2024-06-04, 12:17:53 UTC] {logging_mixin.py:188} INFO - Querying CHDB
[2024-06-04, 12:17:53 UTC] {connection.py:408} WARNING - Failed to connect to localhost:9000
Traceback (most recent call last):
.....
clickhouse_driver.errors.NetworkError: Code: 210. Connection refused (localhost:9000)

日志显示Airflow尝试连接ClickHouse的9000端口(CLI/TCP端口),而非配置的8123 HTTP端口。使用clickhouse_driver.Client的连接代码如下:

class ClickHouseConnection:
    connection = None
    def get_connection(connection_name='my_clickhouse'):
        from clickhouse_driver import Client

        if ClickHouseConnection.connection:
            return connection
        db_props = BaseHook.get_connection(connection_name)
        ClickHouseConnection.connection = Client(db_props.host)
        return ClickHouseConnection.connection, db_props

@logger
def clickhouse_connect():
    print("Connecting to ClickHouse Cluster")
    ch_connection, db_props = ClickHouseConnection.get_connection()
    print(db_props.port)
    print("Connected successfully")
    filename = '/opt/airflow/data/test_csv_file.csv'
    if ch_connection:
        print("Querying CHDB")
        ch_connection.execute(
            f"INSERT INTO maindb.monitor FROM INFILE {filename} FORMAT CSV"
        )

已尝试修改Airflow连接端口为9000并重启容器,但问题未解决,需解决连接失败或强制使用8123端口的问题。

解决方案

核心原因

clickhouse_driver是ClickHouse的原生TCP客户端,默认使用9000端口;而8123是ClickHouse的HTTP接口端口,对应客户端为clickhouse-connect(HTTP客户端)。当前代码仅传入host参数,未指定端口,客户端自动使用默认的9000端口,同时容器网络环境下localhost指向Airflow容器自身,导致连接拒绝。

路径1:修复TCP客户端连接(使用9000端口)

  1. 确保ClickHouse容器暴露9000端口
    启动ClickHouse容器时需映射9000端口,示例命令:

    docker run -p 8123:8123 -p 9000:9000 clickhouse/clickhouse-server
    

    若使用docker-compose,需在ClickHouse服务的ports配置中添加9000:9000。

  2. 修正Airflow连接的Host配置
    Airflow容器内的localhost指向自身,需将连接的Host改为:

    • ClickHouse容器的服务名(docker-compose环境下)
    • WSL的IP地址(可通过ip addr在WSL中查看,通常为172.17.0.1或类似段)
    • ClickHouse容器的直接IP(通过docker inspect <clickhouse-container-id>获取)
  3. 代码中显式传入端口参数
    修改ClickHouseConnection.get_connection方法,传入Airflow连接配置中的端口:

    ClickHouseConnection.connection = Client(
        host=db_props.host,
        port=db_props.port,
        user=db_props.login,  # 若有用户名需添加
        password=db_props.password  # 若有密码需添加
    )
    

路径2:改用HTTP客户端连接(使用8123端口)

若要使用8123端口,需替换为HTTP客户端clickhouse-connect:

  1. 在Airflow容器中安装依赖

    pip install clickhouse-connect
    

    若使用docker-compose,可在Airflow服务的command中添加安装命令,或构建自定义镜像。

  2. 修改连接代码

    class ClickHouseConnection:
        connection = None
        def get_connection(connection_name='my_clickhouse'):
            import clickhouse_connect
    
            if ClickHouseConnection.connection:
                return connection
            db_props = BaseHook.get_connection(connection_name)
            # 创建HTTP连接,默认使用8123端口
            ClickHouseConnection.connection = clickhouse_connect.get_client(
                host=db_props.host,
                port=db_props.port,
                username=db_props.login,
                password=db_props.password
            )
            return ClickHouseConnection.connection, db_props
    
    @logger
    def clickhouse_connect():
        print("Connecting to ClickHouse Cluster")
        ch_connection, db_props = ClickHouseConnection.get_connection()
        print(db_props.port)
        print("Connected successfully")
        filename = '/opt/airflow/data/test_csv_file.csv'
        if ch_connection:
            print("Querying CHDB")
            # 使用文件流完成CSV插入(适配clickhouse-connect的API)
            with open(filename, 'rb') as f:
                ch_connection.insert("maindb.monitor", f, format='CSV')
    

内容的提问来源于stack exchange,提问作者DiploDog

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 21:25:56