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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 19:17:41