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

Apache Airflow 2.8.4镜像中OpenLineage实现失败求助

Airflow中Snowflake、Databricks Operator的OpenLineage实现问题及解决方案

我正在参考Airflow OpenLineage提供商文档,尝试为Snowflake、Databricks等Operator实现OpenLineage,但始终无法得到预期结果。作为Airflow和OpenLineage的新手,我需要相关帮助。以下是我的DAG文件和docker-compose配置文件:

DAG文件内容

from __future__ import annotations

from airflow.providers.databricks.operators.databricks_sql import (
    DatabricksCopyIntoOperator,
    DatabricksSqlOperator,
)

from airflow.providers.snowflake.transfers.copy_into_snowflake import CopyFromExternalStageToSnowflakeOperator
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
import os, boto3
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from datetime import datetime
from airflow.operators.dummy import DummyOperator

SNOWFLAKE_CONN_ID = "my_snowflake_conn"
SNOWFLAKE_STAGE = "kartike"
SNOWFLAKE_SAMPLE_TABLE = "sample_table2"
S3_FILE_PATH = "orders_data_header.csv"
ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID")
DAG_ID = "example_s3_to_snowflake"
DATABRICKS_CONN_ID = "my_databricks_conn"
SQL_ENDPOINT_NAME = "a2fe5ea1499dda95"

# Define params for Run Now Operator
notebook_params = {"Variable": 5}

with DAG(
        DAG_ID,
        start_date=datetime(2021, 1, 1),
        default_args={"snowflake_conn_id": SNOWFLAKE_CONN_ID},
        tags=["example"],
        schedule="@once",
        catchup=False,
) as dag:

    copy_into_table_snowflake = CopyFromExternalStageToSnowflakeOperator(
        task_id="copy_into_table_snowflake",
        snowflake_conn_id=SNOWFLAKE_CONN_ID,
        files=[S3_FILE_PATH],
        table=SNOWFLAKE_SAMPLE_TABLE,
        stage=SNOWFLAKE_STAGE,
        file_format="(type = 'CSV',field_delimiter = ',')",
        pattern=".*[.]csv",
    )


    create_snowflake_table_new_instance = SnowflakeOperator(
        task_id="create_snowflake_table_new_instance",
        snowflake_conn_id=SNOWFLAKE_CONN_ID,
        sql= "INSERT INTO demo_table (CUSTOMER_ID, CUSTOMER_NAME,TYPE) VALUES ('CUST1','NAME1','TYPE1'), ('CUST2','NAME2','TYPE1'), ('CUST3','NAME3','TYPE2');"
    )

    def test_s3_connection():
       #print("S3 configured successfully!")
       aws_access_key_id = 'AKIAEXAMPLE',
       aws_secret_access_key = 'EXAMPLE',
       bucket_name = "arn:aws:s3:::marcl-astrosdk-kk",
       # bucket = str(bucket_name)  # Convert bucket to string if it's not already
       # s3_client = boto3.client('s3', aws_access_key_id=aws_access_key_id, aws_secret_access_key=aws_secret_access_key)
       # objects = s3_client.list_objects_v2(Bucket=bucket_name)
       # if 'Contents' in objects:
       #  print("S3 connection successful!")
       # else:
       #  print("S3 connection failed.")


    AWS_S3 = PythonOperator(
        task_id='AWS_S3',
        python_callable=test_s3_connection,
        dag=dag,
    )

    t0 = DummyOperator(
        task_id='start'
    )


    databricks_task = DatabricksSqlOperator(
        databricks_conn_id=DATABRICKS_CONN_ID,
        http_path = "/sql/1.0/warehouses/a2fe5ea1499dda95",
        task_id="create_and_populate_table",
        sql=[
            "CREATE TABLE IF NOT EXISTS lineage_data.lineagedemo.menu (recipe_id INT, app string, main string,  dessert string)",
            "INSERT INTO lineage_data.lineagedemo.menu (recipe_id, app, main, dessert) VALUES (1,'Ceviche', 'Tacos', 'Flan')",
            "CREATE TABLE IF NOT EXISTS lineage_data.lineagedemo.dinner AS SELECT recipe_id, concat(app, main, dessert) AS full_menu FROM lineage_data.lineagedemo.menu"
            ],
        )

    AWS_S3 >> t0 >> databricks_task >> copy_into_table_snowflake >> create_snowflake_table_new_instance

Docker-Compose配置文件内容

---
x-airflow-common:
  &airflow-common
  image: ${AIRFLOW_IMAGE_NAME:-apache/airflow:2.8.4}
  #image: test
  # build: .
  environment:
    &airflow-common-env
    AIRFLOW__CORE__EXECUTOR: CeleryExecutor
    AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow
    AIRFLOW__CELERY__RESULT_BACKEND: db+postgresql://airflow:airflow@postgres/airflow
    AIRFLOW__CELERY__BROKER_URL: redis://:@redis:6379/0
    AIRFLOW__CORE__FERNET_KEY: ''
    AIRFLOW__CORE__DAGS_ARE_PAUSED_AT_CREATION: 'true'
    AIRFLOW__CORE__LOAD_EXAMPLES: 'true'
    AIRFLOW__SCHEDULER__ENABLE_HEALTH_CHECK: 'true'
    _PIP_ADDITIONAL_REQUIREMENTS: ${_PIP_ADDITIONAL_REQUIREMENTS:-  apache-airflow-providers-databricks}
    AIRFLOW__OPENLINEAGE__TRANSPORT: '{"type": "http", "url": "http://localhost:3000"}'
    AIRFLOW__OPENLINEAGE__NAMESPACE: 'namespaces'
    ##AIRFLOW__OPENLINEAGE__EXTRACTORS: 'plugins.extractors.some_lineage_extractor.MyExtractor'
    AIRFLOW__API__AUTH_BACKENDS: 'airflow.api.auth.backend.basic_auth,airflow.api.auth.backend.session'
#   AIRFLOW__OPENLINEAGE__CONFIG_PATH: '/opt/airflow/plugins/openlineage.yml'
  volumes:
    - ${AIRFLOW_PROJ_DIR:-.}/dags:/opt/airflow/dags
    - ${AIRFLOW_PROJ_DIR:-.}/logs:/opt/airflow/logs
    - ${AIRFLOW_PROJ_DIR:-.}/config:/opt/airflow/config
    - ${AIRFLOW_PROJ_DIR:-.}/plugins:/opt/airflow/plugins
  user: "${AIRFLOW_UID:-50000}:0"
  depends_on:
    &airflow-common-depends-on
    redis:
      condition: service_healthy
    postgres:
      condition: service_healthy

services:
  postgres:
    image: postgres:13
    environment:
      POSTGRES_USER: airflow
      POSTGRES_PASSWORD: airflow
      POSTGRES_DB: airflow
    volumes:
      - postgres-db-volume:/var/lib/postgresql/data
    healthcheck:
      test: ["CMD", "pg_isready", "-U", "airflow"]
      interval: 10s
      retries: 5
      start_period: 5s
    restart: always

  redis:
    image: redis:latest
    expose:
      - 6379
    healthcheck:
      test: ["CMD", "redis-cli", "ping"]
      interval: 10s
      timeout: 30s
      retries: 50
      start_period: 30s
    restart: always

  airflow-webserver:
    <<: *airflow-common
    command: webserver
    ports:
      - "8080:8080"
    healthcheck:
      test: ["CMD", "curl", "--fail", "http://localhost:8080/health"]
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-init:
        condition: service_completed_successfully

  airflow-scheduler:
    <<: *airflow-common
    command: scheduler
    healthcheck:
      test: ["CMD", "curl", "--fail", "http://localhost:8974/health"]
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-init:
        condition: service_completed_successfully

  airflow-worker:
    <<: *airflow-common
    command: celery worker
    healthcheck:
      # yamllint disable rule:line-length
      test:
        - "CMD-SHELL"
        - 'celery --app airflow.providers.celery.executors.celery_executor.app inspect ping -d "celery@$${HOSTNAME}" || celery --app airflow.executors.celery_executor.app inspect ping -d "celery@$${HOSTNAME}"'
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    environment:
      <<: *airflow-common-env
      DUMB_INIT_SETSID: "0"
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-init:
        condition: service_completed_successfully

  airflow-triggerer:
    <<: *airflow-common
    command: triggerer
    healthcheck:
      test: ["CMD-SHELL", 'airflow jobs check --job-type TriggererJob --hostname "$${HOSTNAME}"']
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-init:
        condition: service_completed_successfully

  airflow-init:
    <<: *airflow-common
    entrypoint: /bin/bash
    # yamllint disable rule:line-length
    command:
      - -c
      - |
        if [[ -z "${AIRFLOW_UID}" ]]; then
          echo
          echo -e "\033[1;33mWARNING!!!: AIRFLOW_UID not set!\e[0m"
          echo "If you are on Linux, you SHOULD follow the instructions below to set "
          echo "AIRFLOW_UID environment variable, otherwise files will be owned by root."
          echo "For other operating systems you can get rid of the warning with manually created .env file:"
          echo "    See: https://airflow.apache.org/docs/apache-airflow/stable/howto/docker-compose/index.html#setting-the-right-airflow-user"
          echo
        fi
        one_meg=1048576
        mem_available=$$(($$(getconf _PHYS_PAGES) * $$(getconf PAGE_SIZE) / one_meg))
        cpus_available=$$(grep -cE 'cpu[0-9]+' /proc/stat)
        disk_available=$$(df / | tail -1 | awk '{print $$4}')
        warning_resources="false"
        if (( mem_available < 4000 )) ; then
          echo
          echo -e "\033[1;33mWARNING!!!: Not enough memory available for Docker.\e[0m"
          echo "At least 4GB of memory required. You have $$(numfmt --to iec $$((mem_available * one_meg)))"
          echo
          warning_resources="true"
        fi
        if (( cpus_available < 2 )); then
          echo
          echo -e "\033[1;33mWARNING!!!: Not enough CPUS available for Docker.\e[0m"
          echo "At least 2 CPUs recommended. You have $${cpus_available}"
          echo
          warning_resources="true"
        fi
        if (( disk_available < one_meg * 10 )); then
          echo
          echo -e "\033[1;33mWARNING!!!: Not enough Disk space available for Docker.\e[0m"
          echo "At least 10 GBs recommended. You have $$(numfmt --to iec $$((disk_available * 1024 )))"
          echo
          warning_resources="true"
        fi
        if [[ $${warning_resources} == "true" ]]; then
          echo
          echo -e "\033[1;33mWARNING!!!: You have not enough resources to run Airflow (see above)!\e[0m"
          echo "Please follow the instructions to increase amount of resources available:"
          echo "   https://airflow.apache.org/docs/apache-airflow/stable/howto/docker-compose/index.html#before-you-begin"
          echo
        fi
        mkdir -p /sources/logs /sources/dags /sources/plugins
        chown -R "${AIRFLOW_UID}:0" /sources/{logs,dags,plugins}
        exec /entrypoint airflow version
    # yamllint enable rule:line-length
    environment:
      <<: *airflow-common-env
      _AIRFLOW_DB_MIGRATE: 'true'
      _AIRFLOW_WWW_USER_CREATE: 'true'
      _AIRFLOW_WWW_USER_USERNAME: ${_AIRFLOW_WWW_USER_USERNAME:-airflow}
      _AIRFLOW_WWW_USER_PASSWORD: ${_AIRFLOW_WWW_USER_PASSWORD:-airflow}
      _PIP_ADDITIONAL_REQUIREMENTS: ''
    user: "0:0"
    volumes:
      - ${AIRFLOW_PROJ_DIR:-.}:/sources

  airflow-cli:
    <<: *airflow-common
    profiles:
      - debug
    environment:
      <<: *airflow-common-env
      CONNECTION_CHECK_MAX_COUNT: "0"
    # Workaround for entrypoint issue. See: https://github.com/apache/airflow/issues/16252
    command:
      - bash
      - -c
      - airflow

  # You can enable flower by adding "--profile flower" option e.g. docker-compose --profile flower up
  # or by explicitly targeted on the command line e.g. docker-compose up flower.
  # See: https://docs.docker.com/compose/profiles/
  flower:
    <<: *airflow-common
    command: celery flower
    profiles:
      - flower
    ports:
      - "5555:5555"
    healthcheck:
      test: ["CMD", "curl", "--fail", "http://localhost:5555/"]
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-init:
        condition: service_completed_successfully

volumes:
  postgres-db-volume:

问题解决步骤

  1. 补全依赖安装
    修改docker-compose中_PIP_ADDITIONAL_REQUIREMENTS,添加OpenLineage核心及对应提供商依赖:

    _PIP_ADDITIONAL_REQUIREMENTS: ${_PIP_ADDITIONAL_REQUIREMENTS:-apache-airflow-providers-databricks apache-airflow-providers-openlineage apache-airflow-providers-snowflake openlineage-python}
    
  2. 修正OpenLineage传输地址
    容器内无法直接访问宿主机localhost,替换为Docker宿主机别名或实际IP:

    AIRFLOW__OPENLINEAGE__TRANSPORT: '{"type": "http", "url": "http://host.docker.internal:3000"}'
    

    (host.docker.internal适用于Docker Desktop环境,Linux环境需替换为宿主机实际IP)

  3. 验证Operator版本兼容性

    • 确保apache-airflow-providers-snowflake版本≥5.0.0,apache-airflow-providers-databricks版本≥4.0.0,才能支持OpenLineage自动提取。
    • Databricks侧需开启OpenLineage集成,或在DatabricksSqlOperator中通过openlineage_enabled=True显式启用。
  4. 开启调试日志排查
    添加环境变量开启OpenLineage调试日志:

    AIRFLOW__LOGGING__LOGGING_LEVEL: 'INFO'
    AIRFLOW__OPENLINEAGE__DEBUG: 'true'
    

    查看Airflow worker和scheduler日志,排查传输失败、提取器未加载等问题。

  5. 手动注入Lineage(可选)
    若自动提取失效,可通过OpenLineageHook手动添加血缘信息:

    from openlineage.client.facet import SchemaDatasetFacet, Field
    from airflow.providers.openlineage.plugins.hooks import OpenLineageHook
    
    def add_custom_lineage(**kwargs):
        hook = OpenLineageHook()
        # 注入输入数据集
        hook.emit(
            task_instance=kwargs['ti'],
            inputs=[
                {
                    "namespace": "snowflake
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 03:00:18