1. 依赖
MySqlOperator 的数据库交互通过 MySQLdb 模块来实现, 使用前需要安装相关依赖:
pip install apache-airflow[mysql]
2. 使用
使用 MySqlOperator 执行sql任务的一个简单例子:
1from airflow import DAG 2from airflow.utils.dates import days_ago 3from airflow.operators.mysql_operator import MySqlOperator 4 5default_args = { 6 'owner': 'airflow', 7 'depends_on_past': False, 8 'start_date': days_ago(1), 9 'email': ['j_hao104@163.com'], 10 'email_on_failure': True, 11 'email_on_retry': False, 12} 13 14dag = DAG( 15 'MySqlOperatorExample', 16 default_args=default_args, 17 description='MySqlOperatorExample', 18 schedule_interval="30 18 * * *") 19 20insert_sql = "insert into log SELECT * FROM temp_log" 21 22 23task = MySqlOperator( 24 task_id='select_sql', 25 sql=insert_sql, 26 mysql_conn_id='mysql_conn', 27 autocommit=True, 28 dag=dag)
3. 参数
MySqlOperator 接收几个参数:
sql: 待执行的sql语句;mysql_conn_id: mysql数据库配置ID, Airflow的conn配置有两种配置方式,一是通过os.environ来配置环境变量实现,二是通过web界面配置到代码中,具体的配置方法会在下文描述;parameters: 相当于MySQLdb库的execute方法的第二参数,比如:cur.execute('insert into UserInfo values(%s,%s)',('alex',18));autocommit: 自动执行commit;database: 用于覆盖conn配置中的数据库名称, 这样方便于连接统一个mysql的不同数据库;
4. conn配置

建议conn配置通过web界面来配置,这样不用硬编码到代码中,关于配置中的各个参数:
Conn Id: 对应MySqlOperator中的mysql_conn_id;Host: 数据库IP地址;Schema: 库名, 可以被MySqlOperator中的database重写;Login: 登录用户名;Password: 登录密码;Port: 数据库端口;Extra:MySQLdb.connect的额外参数,包含charset、cursor、ssl、local_infile
其中cursor的值的对应关系为: sscursor —> MySQLdb.cursors.SSCursor; dictcursor —> MySQLdb.cursors.DictCursor; ssdictcursor —> MySQLdb.cursors.SSDictCursor