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

Airflow DAG任务MySQL连接失败问题修复请求

修复Airflow DAG中MySQL连接失败问题

问题场景

我搭建了一个包含三个任务的Airflow DAG:

  • 任务1:抓取https://www.lushusa.com/bath/bath-bombs/的产品数据
  • 任务2:将数据存入MySQL数据库
  • 任务3:查询并展示数据库中的数据

目前任务1执行成功,但任务2因连接问题失败,报错信息如下:

[2023-04-02, 10:25:33 UTC] {standard_task_runner.py:100} ERROR - Failed to execute job 212 for task store_lush_data ('NoneType' object has no attribute 'cursor'; 20486)
[2023-04-02, 10:25:33 UTC] {local_task_job.py:212} INFO - Task exited with return code 1
[2023-04-02, 10:25:33 UTC] {taskinstance.py:2585} INFO - 0 downstream tasks scheduled from follow-on schedule check

原代码如下:

import mysql.connector
from mysql.connector import Error
import pandas as pd
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from bs4 import BeautifulSoup
from airflow.utils.dates import days_ago
import requests

default_args = {
    'owner': 'sehrish',
    'depends_on_past': False,
    'start_date': days_ago(0),
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

dag = DAG(
    dag_id='lush_data_dag',
    default_args=default_args,
    description='A pipeline to scrape and store data from LUSH USA website',
    schedule_interval=timedelta(days=1),
)

def create_server_connection(hostname, username, password):
    connection = None
    try:
        connection = mysql.connector.connect(
            host=hostname,
            user=username,
            passwd=password
        )
        print("MySQL Database connection successful")
    except Error as err:
        print(f"Error: '{err}'")

    return connection

def scrape_lush_data():
    url = 'https://www.lushusa.com/bath/bath-bombs/'
    response = requests.get(url)
    soup = BeautifulSoup(response.text, 'html.parser')
    products = soup.find_all('div', {'class': 'product-tile js-product-tile'})
    data = []
    for product in products:
        name = product.find('span', {'class': 'product-tile__name js-product-tile-name'}).text.strip()
        price = product.find('span', {'class': 'product-tile__price'}).text.strip()
        data.append((name, price))
    return data


def store_lush_data():
    # connection = create_server_connection("localhost", "HHJ", "babyMuhammad1!")
    hostname = "localhost"
    username = "HHJ"
    password = "babyMuhammad1!"
    connection = create_server_connection(hostname, username, password)
    cursor = connection.cursor()
    cursor.execute('USE lush_db')
    cursor.execute('CREATE TABLE IF NOT EXISTS lush_products (id INT IDENTITY(1,1) PRIMARY KEY, name VARCHAR(255) NOT NULL, price VARCHAR(50) NOT NULL);')
    data = scrape_lush_data()
    cursor.executemany('INSERT INTO lush_products (name, price) VALUES (%s, %s);', data)
    connection.commit()
    cursor.close()
    connection.close()

def display_lush_data():
    connection = create_server_connection("localhost", "HHJ", "babyMuhammad1!")
    cursor = connection.cursor()
    cursor.execute('USE lush_db')
    cursor.execute('SELECT * FROM lush_products;')
    rows = cursor.fetchall()
    for row in rows:
        print(row)
    cursor.close()
    connection.close()


t1 = PythonOperator(
    task_id='scrape_lush_data',
    python_callable=scrape_lush_data,
    dag=dag,
)

t2 = PythonOperator(
    task_id='store_lush_data',
    python_callable=store_lush_data,
    dag=dag,
)

t3 = PythonOperator(
    task_id='display_lush_data',
    python_callable=display_lush_data,
    dag=dag,
)

t1 >> t2 >> t3

问题分析与修复方案

1. 空连接未处理引发的错误

原create_server_connection函数在连接失败时仅打印错误,仍返回None,后续调用connection.cursor()必然触发'NoneType' object has no attribute 'cursor'错误。

  • 修复:连接失败时抛出异常,终止任务,避免无效操作。

2. MySQL建表语法错误

原建表语句中id INT IDENTITY(1,1) PRIMARY KEY是SQL Server的自增语法,MySQL需替换为AUTO_INCREMENT。

3. 任务间数据传递方式不合理

原代码在任务2中重复调用任务1的爬取函数,不符合Airflow任务间数据传递规范,应通过XCom传递爬取结果,减少重复请求。

4. 数据库连接优化

连接时直接指定目标数据库,无需单独执行USE lush_db命令。

修改后的完整代码

import mysql.connector
from mysql.connector import Error
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from bs4 import BeautifulSoup
from airflow.utils.dates import days_ago
import requests

default_args = {
    'owner': 'sehrish',
    'depends_on_past': False,
    'start_date': days_ago(0),
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

dag = DAG(
    dag_id='lush_data_dag',
    default_args=default_args,
    description='A pipeline to scrape and store data from LUSH USA website',
    schedule_interval=timedelta(days=1),
)

def create_server_connection(hostname, username, password, db_name):
    connection = None
    try:
        connection = mysql.connector.connect(
            host=hostname,
            user=username,
            passwd=password,
            database=db_name
        )
        print("MySQL Database connection successful")
    except Error as err:
        # 抛出异常,让Airflow捕获并标记任务失败
        raise Exception(f"Database connection failed: {err}")

    return connection

def scrape_lush_data(**context):
    url = 'https://www.lushusa.com/bath/bath-bombs/'
    response = requests.get(url)
    soup = BeautifulSoup(response.text, 'html.parser')
    products = soup.find_all('div', {'class': 'product-tile js-product-tile'})
    data = []
    for product in products:
        name = product.find('span', {'class': 'product-tile__name js-product-tile-name'}).text.strip()
        price = product.find('span', {'class': 'product-tile__price'}).text.strip()
        data.append((name, price))
    # 将爬取结果存入XCom,供下游任务使用
    context['ti'].xcom_push(key='lush_products', value=data)
    return data


def store_lush_data(**context):
    hostname = "localhost"
    username = "HHJ"
    password = "babyMuhammad1!"
    db_name = "lush_db"
    
    # 获取任务1的爬取结果
    data = context['ti'].xcom_pull(key='lush_products', task_ids='scrape_lush_data')
    if not data:
        raise Exception("No data received from scrape_lush_data task")
    
    connection = create_server_connection(hostname, username, password, db_name)
    cursor = connection.cursor()
    
    # 修正MySQL自增语法
    cursor.execute('CREATE TABLE IF NOT EXISTS lush_products (id INT AUTO_INCREMENT PRIMARY KEY, name VARCHAR(255) NOT NULL, price VARCHAR(50) NOT NULL);')
    cursor.executemany('INSERT INTO lush_products (name, price) VALUES (%s, %s);', data)
    
    connection.commit()
    cursor.close()
    connection.close()

def display_lush_data():
    hostname = "localhost"
    username = "HHJ"
    password = "babyMuhammad1!"
    db_name = "lush_db"
    
    connection = create_server_connection(hostname, username, password, db_name)
    cursor = connection.cursor()
    
    cursor.execute('SELECT * FROM lush_products;')
    rows = cursor.fetchall()
    for row in rows:
        print(row)
    
    cursor.close()
    connection.close()


t1 = PythonOperator(
    task_id='scrape_lush_data',
    python_callable=scrape_lush_data,
    provide_context=True,  # 开启上下文传递,支持XCom操作
    dag=dag,
)

t2 = PythonOperator(
    task_id='store_lush_data',
    python_callable=store_lush_data,
    provide_context=True,
    dag=dag,
)

t3 = PythonOperator(
    task_id='display_lush_data',
    python_callable=display_lush_data,
    dag=dag,
)

t1 >> t2 >> t3

额外优化建议

  • 使用Airflow的Connections功能管理数据库凭证,避免代码中明文存储密码:在Airflow UI中创建MySQL连接,然后通过BaseHook.get_connection获取连接信息。
  • 为数据库操作添加事务处理,确保数据一致性。
  • 对爬取结果做异常处理,避免因网页结构变化导致任务失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 20:18:08