如何从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
注意事项
- XCom序列化:Airflow默认支持pandas DataFrame的序列化传递,直接返回DataFrame即可在后续任务中使用。
- 任务依赖补全:原代码缺少任务实例定义,需补充后才能建立依赖关系,示例:
@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
相关产品推荐
相关产品推荐

