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:
问题解决步骤
补全依赖安装
修改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}修正OpenLineage传输地址
容器内无法直接访问宿主机localhost,替换为Docker宿主机别名或实际IP:AIRFLOW__OPENLINEAGE__TRANSPORT: '{"type": "http", "url": "http://host.docker.internal:3000"}'(
host.docker.internal适用于Docker Desktop环境,Linux环境需替换为宿主机实际IP)验证Operator版本兼容性
- 确保
apache-airflow-providers-snowflake版本≥5.0.0,apache-airflow-providers-databricks版本≥4.0.0,才能支持OpenLineage自动提取。 - Databricks侧需开启OpenLineage集成,或在
DatabricksSqlOperator中通过openlineage_enabled=True显式启用。
- 确保
开启调试日志排查
添加环境变量开启OpenLineage调试日志:AIRFLOW__LOGGING__LOGGING_LEVEL: 'INFO' AIRFLOW__OPENLINEAGE__DEBUG: 'true'查看Airflow worker和scheduler日志,排查传输失败、提取器未加载等问题。
手动注入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
相关产品推荐
相关产品推荐

