Airflow 1版本操作Redshift执行SQL与调用存储过程的实现方案
1. 前置依赖说明
Airflow 1.10及以上版本已经在官方contrib模块内置了Redshift操作相关能力,无需自行实现Python Operator即可满足需求。首先确保你的Airflow环境已经安装了对应依赖:pip install 'apache-airflow[postgres,amazon]'
2. 执行普通SQL操作的实现
有两种官方内置Operator可选:
方案1:使用
RedshiftSQLOperator
直接调用Airflow 1.x内置的Redshift专属Operator,示例代码如下:from airflow import DAG from airflow.contrib.operators.redshift_operator import RedshiftSQLOperator from datetime import datetime default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1), } with DAG('redshift_sql_demo', default_args=default_args, schedule_interval=None) as dag: run_custom_sql = RedshiftSQLOperator( task_id='run_custom_sql', # 提前在Airflow后台连接管理中配置Redshift连接,连接类型选择Redshift redshift_conn_id='your_redshift_conn', sql=""" CREATE TABLE IF NOT EXISTS analytics.user_behavior ( user_id BIGINT, visit_time TIMESTAMP, page_url VARCHAR(200) ); """ )方案2:使用
PostgresOperator
如果你的Airflow 1.x版本较旧找不到Redshift专属Operator,可直接用PostgresOperator,因为Redshift完全兼容PostgreSQL通信协议,仅需将连接信息配置为Redshift集群参数即可:from airflow.contrib.operators.postgres_operator import PostgresOperator run_sql = PostgresOperator( task_id='run_sql_via_postgres_operator', # 连接类型选择Postgres,端口填Redshift默认端口5439 postgres_conn_id='your_redshift_conn', sql="SELECT COUNT(*) FROM analytics.user_behavior;" )
3. 触发Redshift存储过程的实现
Airflow完全支持触发Redshift存储过程,直接在上述Operator的sql参数中传入Redshift存储过程调用语句即可,示例如下:
call_stored_proc = RedshiftSQLOperator( task_id='call_redshift_proc', redshift_conn_id='your_redshift_conn', sql="CALL analytics.calculate_daily_uv('2024-01-01');" )
如果需要动态传入存储过程参数,可以直接使用Airflow内置的Jinja模板语法,两个Operator默认支持sql参数的模板渲染。
注意:配置的Redshift连接对应的数据库账号需要拥有对应SQL操作、存储过程的执行权限,否则会触发权限报错。
4. 极端版本兜底方案
如果你的Airflow版本非常老旧,contrib模块没有上述两个Operator,也不需要写Python Operator,可以直接用BashOperator调用psql命令行工具执行操作,前提是Worker节点已经安装psql客户端:
from airflow.operators.bash_operator import BashOperator run_sql_via_cli = BashOperator( task_id='run_sql_via_cli', bash_command='psql "host={{ var.value.redshift_host }} port=5439 dbname={{ var.value.redshift_db }} user={{ var.value.redshift_user }} password={{ var.value.redshift_pwd }}" -c "CALL analytics.calculate_daily_uv(\'{{ ds }}\');"' )
内容的提问来源于stack exchange,提问作者MiepMiep

