如何在Apache Airflow中将AWS S3 Bucket数据导入MariaDB表?
实现S3键导入本地Docker MariaDB的Airflow方案
核心思路
通过Airflow现有Operator组合实现:先列出S3桶内的所有对象键,再批量插入到MariaDB指定表中,无需自定义Operator即可完成需求。
前置准备
- 配置Airflow连接
- 在Airflow UI的「Admin > Connections」中添加:
- AWS连接(ID建议设为
aws_default):填入拥有S3 ListObjects权限的Access Key/Secret Key - MariaDB连接(ID建议设为
mariadb_default):- 主机:如果MariaDB和Airflow同属Docker网络,填容器名;否则填宿主机IP
- 端口:默认3306
- 数据库名、用户名、密码按实际配置填写
- AWS连接(ID建议设为
- 在Airflow UI的「Admin > Connections」中添加:
- 创建MariaDB目标表
在本地MariaDB中执行建表语句(可根据需求调整字段):CREATE TABLE IF NOT EXISTS s3_keys ( id INT AUTO_INCREMENT PRIMARY KEY, key_name VARCHAR(255) NOT NULL UNIQUE, imported_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP );
DAG代码实现
from airflow import DAG from airflow.providers.amazon.aws.operators.s3 import S3ListOperator from airflow.operators.python import PythonOperator from airflow.providers.mysql.hooks.mysql import MySqlHook from airflow.utils.dates import days_ago default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': days_ago(1), 'retries': 1, } with DAG( 's3_keys_to_mariadb', default_args=default_args, description='Import S3 object keys to local MariaDB table', schedule_interval=None, # 按需设置调度,例如@daily catchup=False, ) as dag: # 第一步:列出S3桶内所有对象键 list_s3_objects = S3ListOperator( task_id='list_s3_objects', aws_conn_id='aws_default', bucket_name='your-target-bucket', # 可选:添加前缀过滤特定路径下的对象,例如prefix='data/' # delimiter='/', # 可选:排除文件夹路径 ) # 第二步:批量插入键到MariaDB def insert_keys_to_mariadb(**context): # 从XCom获取上一步的S3键列表 s3_keys = context['ti'].xcom_pull(task_ids='list_s3_objects') # 过滤掉文件夹路径(如果启用了delimiter则可省略) valid_keys = [key for key in s3_keys if not key.endswith('/')] # 建立MariaDB连接 mysql_hook = MySqlHook(mysql_conn_id='mariadb_default') conn = mysql_hook.get_conn() cursor = conn.cursor() # 批量插入SQL(用INSERT IGNORE避免重复键报错) insert_sql = "INSERT IGNORE INTO s3_keys (key_name) VALUES (%s)" values = [(key,) for key in valid_keys] try: cursor.executemany(insert_sql, values) conn.commit() print(f"成功插入 {len(values)} 条S3键记录") except Exception as e: conn.rollback() raise e finally: cursor.close() conn.close() insert_task = PythonOperator( task_id='insert_keys_to_mariadb', python_callable=insert_keys_to_mariadb, provide_context=True, ) # 设置任务依赖 list_s3_objects >> insert_task
关键说明
- 批量插入优化:使用
executemany而非逐条插入,大幅提升大量键的导入效率 - 去重处理:
INSERT IGNORE会跳过已存在的键,若需要更新已有记录,可改用INSERT ... ON DUPLICATE KEY UPDATE imported_at = CURRENT_TIMESTAMP - 权限控制:确保AWS角色拥有目标S3桶的
ListBucket权限,MariaDB用户拥有目标表的INSERT权限 - Docker网络适配:若Airflow和MariaDB均为Docker部署,需将两者加入同一自定义网络,避免连接失败
内容的提问来源于stack exchange,提问作者Logan McNatt
相关产品推荐
相关产品推荐

