使用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
相关产品推荐
相关产品推荐

