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

基于Docker Compose的AirFlow连接SQL Server与MongoDB失败求助

AirFlow连接SQL Server与MongoDB的Docker部署方案

这种部署完全可行,你现有配置的核心问题是缺少MongoDB提供者依赖,且docker-compose.yml未配置AirFlow运行必需的元数据库与调度服务。以下是修正后的完整配置方案:

一、修正后的Dockerfile

需要同时安装SQL Server和MongoDB的AirFlow提供者,以及SQL Server连接依赖的ODBC驱动:

FROM apache/airflow:2.2.0

USER root
# 安装适配AirFlow 2.2.0的SQL Server、MongoDB提供者
RUN pip install apache-airflow-providers-microsoft-mssql==3.2.0 apache-airflow-providers-mongo==2.3.3
# 安装SQL Server ODBC驱动(连接必需)
RUN apt-get update && apt-get install -y unixodbc unixodbc-dev
RUN curl https://packages.microsoft.com/keys/microsoft.asc | apt-key add -
RUN curl https://packages.microsoft.com/config/debian/10/prod.list > /etc/apt/sources.list.d/mssql-release.list
RUN apt-get update && ACCEPT_EULA=Y apt-get install -y msodbcsql17

USER airflow

二、完整的docker-compose.yml

LocalExecutor模式需要依赖PostgreSQL作为元数据库,同时要启动调度器和初始化服务:

version: '3.8'

services:
  postgres:
    image: postgres:13
    environment:
      - POSTGRES_USER=airflow
      - POSTGRES_PASSWORD=airflow
      - POSTGRES_DB=airflow
    volumes:
      - postgres-db-volume:/var/lib/postgresql/data

  airflow-webserver:
    build:
      context: ./
      dockerfile: Dockerfile
    depends_on:
      - postgres
      - airflow-init
    ports:
      - "9001:8080"
    environment:
      - AIRFLOW__CORE__EXECUTOR=LocalExecutor
      - AIRFLOW__CORE__SQL_ALCHEMY_CONN=postgresql+psycopg2://airflow:airflow@postgres/airflow
      - AIRFLOW__WEBSERVER__RBAC=true
      - AIRFLOW__CORE__LOAD_EXAMPLES=false
    volumes:
      - ./dags:/opt/airflow/dags
      - ./logs:/opt/airflow/logs
      - ./plugins:/opt/airflow/plugins
    command: webserver

  airflow-scheduler:
    build:
      context: ./
      dockerfile: Dockerfile
    depends_on:
      - postgres
      - airflow-init
    environment:
      - AIRFLOW__CORE__EXECUTOR=LocalExecutor
      - AIRFLOW__CORE__SQL_ALCHEMY_CONN=postgresql+psycopg2://airflow:airflow@postgres/airflow
      - AIRFLOW__CORE__LOAD_EXAMPLES=false
    volumes:
      - ./dags:/opt/airflow/dags
      - ./logs:/opt/airflow/logs
      - ./plugins:/opt/airflow/plugins
    command: scheduler

  airflow-init:
    build:
      context: ./
      dockerfile: Dockerfile
    environment:
      - AIRFLOW__CORE__EXECUTOR=LocalExecutor
      - AIRFLOW__CORE__SQL_ALCHEMY_CONN=postgresql+psycopg2://airflow:airflow@postgres/airflow
      - AIRFLOW__CORE__LOAD_EXAMPLES=false
    command: version
    volumes:
      - ./dags:/opt/airflow/dags
      - ./logs:/opt/airflow/logs
      - ./plugins:/opt/airflow/plugins

volumes:
  postgres-db-volume:

三、AirFlow连接配置

启动容器后,登录AirFlow UI(http://localhost:9001),添加两个连接:

  • SQL Server连接:

    • 连接ID:mssql_default(可自定义)
    • 连接类型:Microsoft SQL Server
    • 主机:SQL Server的IP/域名
    • 登录:SQL Server用户名
    • 密码:SQL Server密码
    • 端口:默认1433
    • 数据库:目标数据库名
  • MongoDB连接:

    • 连接ID:mongodb_default(可自定义)
    • 连接类型:MongoDB
    • 主机:MongoDB的IP/域名
    • 登录:MongoDB用户名(若开启认证)
    • 密码:MongoDB密码(若开启认证)
    • 端口:默认27017
    • 额外:填写database=目标数据库名(需指定时)

四、数据同步示例DAG

在./dags目录下创建sqlserver_to_mongodb.py,实现从SQL Server同步数据到MongoDB:

from airflow import DAG
from airflow.providers.microsoft.mssql.hooks.mssql import MsSqlHook
from airflow.providers.mongo.hooks.mongo import MongoHook
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

def sync_data():
    # 从SQL Server读取数据
    mssql_hook = MsSqlHook(mssql_conn_id='mssql_default')
    df = mssql_hook.get_pandas_df("SELECT * FROM your_source_table")
    
    # 写入MongoDB
    mongo_hook = MongoHook(mongo_conn_id='mongodb_default')
    client = mongo_hook.get_conn()
    db = client['your_mongo_db']
    collection = db['your_mongo_collection']
    collection.insert_many(df.to_dict('records'))

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    'sqlserver_mongodb_sync',
    default_args=default_args,
    description='Sync data from SQL Server to MongoDB',
    schedule_interval=timedelta(days=1),
    catchup=False,
) as dag:

    sync_task = PythonOperator(
        task_id='sync_sqlserver_to_mongodb',
        python_callable=sync_data
    )

    sync_task

内容的提问来源于stack exchange,提问作者Diwas Poudel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 04:35:57