使用Airflow MySQLHook连接MariaDB查询超时问题求助
解决Airflow MySQLHook连接MariaDB查询超时问题
问题背景
查询一张500万行、15列的comment_table表,需获取modified_date(timestamp类型)指定日期的所有行,运行2小时后超时,错误信息:
(2013, 'Lost connection to server during query') unable to rollback;1009)
最初使用的查询语句:
select * from comment_table where DATE(modified_date) = '2023-06-01';
尝试按6小时间隔拆分查询后仍超时:
SELECT * FROM comment_table WHERE modified_date between str_to_date('2023-06-01 00:00:00', '%Y-%m-%d %H:%i:%s') AND str_to_date('2023-06-01 06:00:00', '%Y-%m-%d %H:%i:%s');
使用的核心代码:
hook = MySqlHook("mariadb_conn_id") with hook.get_conn().cursor() as cursor: query = 'select * from comment_table.....' df_results = pd.read_sql(query, hook.get_conn())
解决方案
1. 修复查询语句,让索引生效
对modified_date使用DATE()函数会导致索引失效,改用直接的时间范围匹配,同时确保字段有索引:
- 优化后的查询语句:
SELECT * FROM comment_table WHERE modified_date >= '2023-06-01 00:00:00' AND modified_date < '2023-06-02 00:00:00'; - 给
modified_date创建索引(如果还没建):CREATE INDEX idx_comment_modified_date ON comment_table(modified_date);
2. 分批次读取数据
不要一次性加载所有结果,按更小的时间窗口分批拉取,避免内存和连接超时:
import datetime import pandas as pd from airflow.providers.mysql.hooks.mysql import MySqlHook hook = MySqlHook("mariadb_conn_id") conn = hook.get_conn() batch_window = datetime.timedelta(minutes=30) # 按30分钟为一个批次 start_dt = datetime.datetime(2023, 6, 1, 0, 0, 0) end_dt = datetime.datetime(2023, 6, 2, 0, 0, 0) all_batches = [] current_dt = start_dt while current_dt < end_dt: next_dt = current_dt + batch_window query = f""" SELECT * FROM comment_table WHERE modified_date >= '{current_dt.strftime('%Y-%m-%d %H:%M:%S')}' AND modified_date < '{next_dt.strftime('%Y-%m-%d %H:%M:%S')}' """ df_batch = pd.read_sql(query, conn) all_batches.append(df_batch) current_dt = next_dt # 合并所有批次数据 df_results = pd.concat(all_batches, ignore_index=True) conn.close()
3. 复用连接并调整超时参数
原代码两次调用hook.get_conn()会创建两个连接,改成复用同一个连接;同时增加连接超时配置:
from sqlalchemy import create_engine from airflow.providers.mysql.hooks.mysql import MySqlHook hook = MySqlHook("mariadb_conn_id") # 获取连接URL并添加超时参数 conn_uri = hook.get_uri().replace("mysql://", "mysql+pymysql://") + "?connect_timeout=3600&read_timeout=3600" engine = create_engine(conn_uri) query = """ SELECT * FROM comment_table WHERE modified_date >= '2023-06-01 00:00:00' AND modified_date < '2023-06-02 00:00:00' """ df_results = pd.read_sql(query, engine)
4. 调整Airflow连接配置
在Airflow的MariaDB连接(mariadb_conn_id)的Extra字段中添加超时参数:
{"connect_timeout": 3600, "read_timeout": 3600}
内容的提问来源于stack exchange,提问作者KristiLuna
相关产品推荐
相关产品推荐

