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

如何在Apache Airflow中将AWS S3 Bucket数据导入MariaDB表?

实现S3键导入本地Docker MariaDB的Airflow方案

核心思路

通过Airflow现有Operator组合实现:先列出S3桶内的所有对象键,再批量插入到MariaDB指定表中,无需自定义Operator即可完成需求。

前置准备

  1. 配置Airflow连接
    • 在Airflow UI的「Admin > Connections」中添加:
      • AWS连接(ID建议设为aws_default):填入拥有S3 ListObjects权限的Access Key/Secret Key
      • MariaDB连接(ID建议设为mariadb_default):
        • 主机:如果MariaDB和Airflow同属Docker网络,填容器名;否则填宿主机IP
        • 端口:默认3306
        • 数据库名、用户名、密码按实际配置填写
  2. 创建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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 22:35:16