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

如何从Airflow的ClickHouseHook获取pandas.DataFrame?

问题:如何在Airflow DAG中从ClickHouse获取带列名的pandas.DataFrame?

当前Airflow DAG代码可正常运行,但get_table_from_db(sql_query)返回的结果仅为无列名的二维值列表,无法直接得到带列名的pandas.DataFrame,需调整代码实现需求。

解决方案

方法1:利用ClickHouseHook内置方法(推荐)

如果使用的airflow_clickhouse_plugin版本支持,直接调用get_pandas_df方法即可一键生成带列名的DataFrame:

修改get_table_from_db任务代码:

@task
def get_table_from_db(sql_query):
    # 通过hook直接获取带列名的pandas DataFrame
    df = ch_hook.get_pandas_df(sql_query)
    return df

方法2:自定义转换函数(兼容所有版本)

若hook无内置DataFrame生成方法,可通过游标(cursor)的描述信息提取列名,手动构建DataFrame:

修改get_table_from_db任务代码:

@task
def get_table_from_db(sql_query):
    with ch_hook.get_conn() as conn:
        with conn.cursor() as cursor:
            cursor.execute(sql_query)
            # 从游标描述中提取列名
            column_names = [col[0] for col in cursor.description]
            # 获取查询结果数据
            data_rows = cursor.fetchall()
            # 构建带列名的DataFrame
            df = pd.DataFrame(data_rows, columns=column_names)
    return df

注意事项

  1. XCom序列化:Airflow默认支持pandas DataFrame的序列化传递,直接返回DataFrame即可在后续任务中使用。
  2. 任务依赖补全:原代码缺少任务实例定义,需补充后才能建立依赖关系,示例:
@dag(default_args=default_args, schedule_interval=schedule_interval, catchup=False, concurrency=3)
def test_dag_4():
    # 替换为实际查询语句
    sql_query = "SELECT id, name FROM your_target_table"
    
    @task
    def get_table_from_db(sql_query):
        # 上述方法的实现代码
        ...
    
    @task
    def send_msg(bot_token: str, chat_id: str, df: pd.DataFrame):
        # 将DataFrame转为适合发送的文本格式(如Markdown表格)
        message = df.to_markdown()
        url = f'https://api.telegram.org/bot{bot_token}/sendMessage?chat_id={chat_id}&text={message}'
        client = httpx.Client()
        return client.post(url)
    
    # 创建任务实例并设置依赖
    df_task = get_table_from_db(sql_query)
    send_msg_task = send_msg(bot_token, chat_id, df_task)
    df_task >> send_msg_task

st_dag = test_dag_4()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 00:12:06