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

使用Airflow MySqlOperator执行Insert查询时出现语法错误求助

解决MySqlOperator插入CSV数据时的语法错误问题

问题根源

你当前的代码把SQL语句+参数组成的元组传给了MySqlOperator的sql参数,但MySqlOperator的sql参数仅接受纯SQL字符串,它不会自动解析并替换%s占位符。MySQL服务器收到insert into test1 values (%s,%s)这个SQL时,会把%s当成无效语法,直接抛出1064错误。

解决方案

下面提供两种可行的解决思路,优先推荐第一种(批量导入,效率更高):


方案1:使用MySQL的LOAD DATA INFILE批量导入CSV

这是处理CSV导入MySQL最高效的方式,不需要逐行循环,直接让MySQL读取文件完成导入:

from airflow.providers.mysql.operators.mysql import MySqlOperator
import os
from datetime import datetime

path = "/your/csv/directory"
for filename in os.listdir(path):
    if filename.endswith('.csv'):
        # 构建LOAD DATA的SQL语句,根据你的CSV格式调整参数(比如分隔符、是否有表头)
        sql = f"""
        LOAD DATA LOCAL INFILE '{os.path.join(path, filename)}'
        INTO TABLE test1
        FIELDS TERMINATED BY ',' 
        LINES TERMINATED BY '\n'
        -- 如果CSV有表头,加上下面这行跳过第一行
        -- IGNORE 1 ROWS
        (col1, col2); -- 替换成你test1表的实际列名
        """
        mysql_op = MySqlOperator(
            task_id=f'import_{filename}',
            sql=sql,
            mysql_conn_id='hack5_id',
            owner='hack5',
            dag=dag,
            # 必须开启local_infile才能使用LOAD DATA LOCAL INFILE
            params={'local_infile': True}
        )
        mysql_op.run(start_date=datetime.now(), end_date=datetime(2018, 5, 21))

方案2:改用PythonOperator逐行插入(适合小批量数据)

如果必须逐行处理数据,可以用PythonOperator,在Python函数中直接使用MySQL连接执行带参数的插入,这样能正确解析%s占位符:

from airflow.operators.python import PythonOperator
from airflow.providers.mysql.hooks.mysql import MySqlHook
import csv
import os
from datetime import datetime

def import_csv_to_mysql(**context):
    path = "/your/csv/directory"
    mysql_hook = MySqlHook(mysql_conn_id='hack5_id')
    conn = mysql_hook.get_conn()
    cursor = conn.cursor()
    
    for filename in os.listdir(path):
        if filename.endswith('.csv'):
            file_path = os.path.join(path, filename)
            with open(file_path, 'r') as f:
                csv_data = csv.reader(f)
                # 如果CSV有表头,跳过第一行
                # next(csv_data)
                for row in csv_data:
                    # 这里用参数绑定,cursor.execute会自动替换%s
                    cursor.execute("INSERT INTO test1 VALUES (%s, %s)", row)
    conn.commit()
    cursor.close()
    conn.close()

# 定义PythonOperator
python_op = PythonOperator(
    task_id='import_csv_task',
    python_callable=import_csv_to_mysql,
    owner='hack5',
    dag=dag
)
python_op.run(start_date=datetime.now(), end_date=datetime(2018, 5, 21))

注意事项

  • 如果你使用LOAD DATA INFILE,需要确保MySQL服务器开启了local_infile参数,并且Airflow的MySQL连接配置中允许使用该参数。
  • 如果CSV文件包含表头,记得在SQL或代码中跳过第一行,避免把表头插入到int类型的列中导致类型错误。
  • 批量导入的性能远高于逐行插入,建议优先选择方案1。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:57:35