ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

OpenMetadata Airflow Pipeline 连接器配置实战:REST API 与数据库双通道元数据采集

2026/9/15 18:08:59 拓冰建站 浏览量
OpenMetadata Airflow Pipeline 连接器配置实战:REST API 与数据库双通道元数据采集 OpenMetadata Airflow Pipeline 连接器配置实战REST API 与数据库双通道元数据采集【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata本文以 OpenMetadata 的 Airflow Pipeline 连接器配置文档 为核心系统讲解如何将 Apache Airflow 的 DAG 拓扑、任务结构、调度与运行状态接入 OpenMetadata 元数据体系。你将掌握三种元数据提取策略的取舍、REST API 连接下四种认证方式Basic Auth、Access Token、GCP Service Account、MWAA的配置细节以及 MySQL / Postgres / MSSQL / SQLite 四种底层数据库直连的完整参数清单并理解连接器底层源码的实现原理与测试连接诊断逻辑。三种 Airflow 元数据提取策略在 OpenMetadata 中从 Airflow 提取元数据有三种互补的策略本连接器文档对应的是第一种策略工作方式适用场景1. Airflow Connector在 OpenMetadata 中配置 Pipeline 服务连接通过访问 Airflow 底层元数据库或 REST API 读取 DAG 元数据本篇文章的主体可在 OpenMetadata UI 中直接配置2. Airflow Lineage Backend在 Airflow 实例内安装 OpenMetadata Lineage Backend 插件随任务执行实时上报血缘需要表级、列级血缘的场景3. Airflow Lineage Operator在 DAG 中直接使用 Lineage Operator 发送元数据希望由 DAG 自身控制上报逻辑的场景从 OpenMetadata UI 中可以直接使用策略 1Airflow Connector。需要特别强调的是Lineage Backend 与 Lineage Operator 并不在本连接器配置范围内它们是独立安装在 Airflow 侧的组件通过发送 OpenLineage 事件血缘边会在 OpenMetadata 中自动出现。基础连接配置Connection Details在 OpenMetadata UI 中创建 Airflow Pipeline 服务时首先需要填写以下基础字段其数据模型定义在 airflowConnection.json。Host and PortPipeline Service 管理 URI必须使用scheme://hostname:port的 URI 格式例如http://localhost:8080、http://host.docker.internal:8080。这是 Schema 中唯一必填的字段之一required: [hostPort, connection]声明为 URI 格式。Number Of Status每次 ingestion 运行时读取的历史任务状态数量默认值为 10即默认获取并更新最近 10 次运行的任务状态。该字段在 Schema 中定义为numberOfStatus类型为 integer默认值10。如果你的 DAG 运行频繁且需要更长的状态回溯窗口可以适当调大此值。Metadata Database Connection此处的connection字段是连接器最核心的配置采用oneOf联合类型可选以下五种连接方式定义见 airflowConnection.jsonAirflow REST API 连接airflowRestApiConnection.json通过 HTTP/HTTPS 调用 Airflow Web Server无需直接访问元数据库Backend ConnectionbackendConnection.json用于从运行在 OpenMetadata 实例中的 DAG 直接提取元数据例如 GCS Composer 场景MySQL 连接直连 Airflow 的 MySQL 元数据库Postgres 连接直连 Airflow 的 Postgres 元数据库SQLite 连接直连 Airflow 的 SQLite 元数据库。Schema 还提供了可选的pipelineFilterPattern字段用于通过正则表达式排除不需要采集的 pipeline。从源码实现看connection.pyAirflowConnection._get_client()会根据connection.connection的实际类型分派若是AirflowRestApiConnection则构建AirflowApiClientREST 通道否则通过singledispatch将 MySQL / Postgres / SQLite 配置分别委托给对应的数据库连接器MySQLConnection、PostgresConnection、SQLiteConnection来创建 SQLAlchemy Engine。Airflow REST API 连接REST API 连接通过 HTTP/HTTPS 调用 Airflow Web Server不需要直接访问 Airflow 元数据库。因此它是托管部署Astronomer、GCP Cloud Composer、MWAA以及任何无法或不愿开放数据库直连的自建 Airflow 的正确选择。其配置模型定义在 airflowRestApiConnection.json。注意REST API 连接只获取 DAG 拓扑、任务结构、调度和运行状态不捕获血缘。如需表级和列级血缘必须在 Airflow 中单独安装 OpenMetadata Lineage Backend策略 2或在 DAG 中使用 Lineage Operator策略 3。一旦这些组件发出 OpenLineage 事件血缘边会自动出现在 OpenMetadata 中。不同部署形态下的 Host URL 格式部署形态示例 Host and Port URL自建 / Dockeringestion 运行在宿主机http://localhost:8080自建 / Dockeringestion 运行在 Docker 容器内http://host.docker.internal:8080Google Cloud Composerhttps://ko82752sdo9f7zjf811c682mw1e5uuc9-dot-us-east1.composer.googleusercontent.comAstronomerhttps://cmn4c1zax823t00qf36gnlquw.ay.astronomer.run/v13jlquw/Amazon MWAAhttps://a1234awd1-5324-6f89-9523-1sq41234adqa.c2.airflow.eu-north-1.on.awsCloud Composer在 GCP Console →Composer → Environments → Open Airflow UI中查看 Web Server URL复制基础 URL去掉尾部路径。Astronomer在 Astronomer UI →Deployments → Open Airflow中查看部署 URL注意不要包含尾部斜杠。何时使用 REST API 而非数据库连接使用 REST API 连接的场景使用 Astronomer数据库不可访问使用 Cloud Composer 或 MWAA数据库不可访问或不实际运行 Airflow 3.x无法直接访问底层的 MySQL / Postgres / SQLite 元数据库。使用数据库连接下文 MySQL / Postgres / SQLite 小节的场景自建 Airflow 且可以直接访问元数据库希望直接从数据库中读取原始 task-instance 数据而非通过 API使用 Backend Connection 策略Airflow 插件 / Lineage Backend 方案。认证方式Authentication ConfigurationauthConfig字段为oneOf联合类型支持四种认证方式见 airflowRestApiConnection.jsonBasic Auth输入用户名和密码。Airflow 3.x 在启动时自动交换短期 JWTAirflow 2.x 直接使用 HTTP Basic 认证Access Token粘贴在 Airflow 中生成的静态 Bearer TokenGCP Service Account推荐用于Google Cloud Composer。通过google-auth在运行时获取并自动刷新 GCP OAuth2 TokenToken 不会在运行中途过期MWAA Configuration用于通过 AWS 凭证认证 Amazon Managed Workflows for Apache AirflowMWAA。认证方式快速参考部署形态推荐认证方式自建 Airflow 2.x 或 3.xBasic AuthAstronomerAccess TokenDeployment API TokenGoogle Cloud ComposerGCP Service Account任何已有预生成 Bearer Token 的部署Access TokenUsername 与 PasswordBasic Auth 的用户名username必须具备调用 Airflow REST API 的权限。在 Airflow 3.x 下该用户名会触发自动 JWT 交换POST /auth/token在 Airflow 2.x 下直接使用 HTTP Basic 认证。password即对应的 Basic Auth 密码。源码佐证在 auth.py 中try_exchange_jwt()会向{host}/auth/token发送 POST 请求换取 JWTbuild_basic_auth_callback()返回的回调函数在每次刷新周期都会重新调用 JWT 交换——成功则返回Bearer jwt刷新间隔 25 分钟远小于 Airflow 约 30–60 分钟的 JWT TTL失败说明是 Airflow 2.x则回退为Basic base64(user:pass)TTL 7 天。这意味着长时运行的 ingestion 不会因为 Token 过期而收到 401。TokenAccess Token 认证静态 Bearer Tokentoken粘贴后每次请求都会以Authorization: Bearer token头发送。适用于已在 Airflow 部署中生成长期有效 API Token 的场景。源码中build_access_token_callback()返回一个无过期时间的静态 Token 回调auth.py。生成 Astronomer Deployment Token打开 Astronomer UI导航到Deployments选择你的部署进入API Keys或Tokens具体标签名视 Astronomer 版本而定点击Add API Key/Generate Token命名如openmetadata-ingestion并复制 Token 值将 Token 粘贴到上方Token字段。自建 Airflow 也可以通过 Airflow UI 的Admin → Users或 Airflow CLI 生成 API Token。MWAA ConfigurationMWAA 认证使用 AWS 凭证认证 Amazon MWAA 环境需要提供MWAA Environment Name和AWS Configuration两组信息其 Schema 定义在 mwaaAuthConfig.json。配置字段清单MWAA Environment Name要连接的 Amazon MWAA 环境名称AWS RegionMWAA 环境所在的 AWS 区域AWS Access Key ID用于认证 AWS 的访问密钥 IDAWS Secret Access Key与该访问密钥关联的密钥AWS Session Token可选使用临时 AWS 凭证时必填Assume Role ARN可选跨账户访问时扮演的 IAM 角色 ARNAssume Role Session Name可选扮演角色时的会话名称Endpoint URL可选AWS 兼容服务MinIO、LocalStack的自定义端点 URL。源码佐证MWAA 通道的实现不走 HTTP Token 管理而是直接使用 AWS SDK 的invoke_rest_api方法调用 Airflow REST API见 mwaa.py。MWAAClient内部通过AWSClient获取mwaa_client_invoke_rest_api()将Name环境名、Path如/dags、Method、Body、QueryParameters组装后调用invoke_rest_api。在 client.py 中当检测到MwaaAuthentication时直接构建MWAAClient不再创建 TrackedREST HTTP 客户端。MWAA 环境的连通性测试通过实际调用一次GET /dags?limit1来完成因为 MWAA 不暴露 version 端点见test_get_version()。GCP CredentialsCloud Composer 认证GCP 凭证用于获取短期 OAuth2 Token 以认证 Google Cloud Composer。Token 过期时自动刷新ingestion 运行不会因 Token 过期而中断。支持全部四种 GCP 认证类型GCP Credentials Values直接粘贴 Service Account JSON 字段project ID、client email、private key 等GCP Credentials Path提供 ingestion 主机上 Service Account JSON 密钥文件的路径GCP External AccountWorkload Identity Federation用于 GKE 或其他 workload identity 场景GCP ADCApplication Default Credentials使用环境中已有的凭证如通过gcloud auth application-default login或 GCE metadata server 获取。此外还可以通过gcpImpersonateServiceAccount配置 Service Account 模拟impersonation。源码佐证在 auth.py 的build_gcp_token_callback()中首先调用set_google_credentials()设置凭证若配置了gcpImpersonateServiceAccount则通过get_gcp_impersonate_credentials()获取模拟凭证否则调用google.auth.default()加载默认凭证随后以cloud-platformscope 刷新 Token并使用凭证的expiry时间驱动自动刷新。查找 Cloud Composer Airflow URL在 GCP Console 中进入Composer → Environments选择你的环境并点击Open Airflow UI复制基础 URL如https://hash-dot-region.composer.googleusercontent.com填入Host and Port字段。GCP 凭证类型选择指南凭证类型适用场景GCP Credentials Valuesingestion 运行在 GCP 外部本地、自建环境直接粘贴 Service Account JSON 字段GCP Credentials Pathingestion 主机上已存在已知路径的 Service Account JSON 密钥文件GCP ADCApplication Default Credentialsingestion 运行在带附加 Service Account 的 GCE VM 或 GKE Pod 上使用 GCE metadata server 或gcloud auth application-default loginGCP External AccountWorkload Identity Federationingestion 运行在使用 Workload Identity 的 GKE 上或使用联合身份的非 GCP 系统如 AWS → GCPAPI VersionAirflow REST API 版本取值如下auto默认OpenMetadata 先尝试v2Airflow 3.x失败后回退到v1Airflow 2.xv1强制使用 Airflow 2.x APIv2强制使用 Airflow 3.x API。源码佐证版本探测逻辑在 client.py 的api_version属性与_detect_api_version()中实现——auto模式下依次请求/{version}/version端点先试v2再试v1401/403 会直接抛出凭证问题其他 HTTP 错误继续尝试下一个版本全部失败则默认v1。MWAA 客户端则内部固定为v1版本管理由 MWAA 处理。此外v1/v2 的字段差异也被封装_date_field在 v2 返回logical_date、v1 返回execution_date调度字段 v2 读取timetable_summary、v1 读取schedule_interval见build_dag_details()。Verify SSL是否在连接 Airflow REST API 时校验 SSL 证书。仅在开发环境使用自签名证书时才设置为false。Schema 中该字段默认值为true见 airflowRestApiConnection.json。在 client.py 中verify_ssl会被透传到ClientConfig与build_basic_auth_callback()同时 REST 客户端配置了retry_codes[503, 504]以应对网关类临时错误。MySQL 连接数据库直连如果 Airflow 使用 MySQL 作为元数据库需要填写以下信息Username Password具备连接数据库权限的凭证只需只读权限Host and PortMySQL 服务地址格式为hostname:port如localhost:3306、host.docker.internal:3306Database Schema包含 Airflow 表的 MySQL schema。SSL 相关字段SSL CAsslCASSL CA 文件的路径文件必须位于 ingestion 进程本地SSL CertificatesslCertSSL 客户端证书文件的路径SSL KeysslKeySSL 密钥文件的路径。Postgres 连接数据库直连如果 Airflow 使用 Postgres 作为元数据库需要填写Username Password具备连接数据库权限的凭证只需只读权限Host and PortPostgres 服务地址如localhost:5432或host.docker.internal:5432Database包含 Airflow 表的 Postgres 数据库。SSL Mode连接 Postgres 的 SSL 模式例如prefer、verify-ca等。文档特别说明其余属性可以忽略因为不会采集任何数据库策略标签此处的 Postgres 连接仅用于 Airflow 元数据读取。MSSQL 连接数据库直连如果 Airflow 使用 MSSQL 作为元数据库需要填写Username具备连接数据库权限的凭证只需只读权限Host and PortMSSQL 服务地址如localhost:1433或host.docker.internal:1433。Auth Config认证配置MSSQL 支持三种认证类型authTypeBasic Auth使用密码认证IAM 认证用于连接 AWS 相关服务Azure 认证用于连接 Azure 相关服务。注意如果使用 IAM 认证需要在 Connection Arguments 中添加ssl: {ssl-mode: allow}。Basic AuthPassword连接 MySQL 的密码注意文档原文如此MSSQL Basic Auth 字段沿用此命名。IAM Auth ConfigAWS 凭证AWS Access Key IDAWS 安全凭证的一部分。访问密钥由两部分组成访问密钥 ID如AKIAIOSFODNN7EXAMPLE和秘密访问密钥如wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY两者必须同时使用才能完成认证AWS Secret Access Key秘密访问密钥AWS Region目标服务所在的 AWS 区域。这是配置连接时唯一必填的参数其余 AWS 配置可通过 SDK 的凭证链自动解析AWS Session Token使用临时凭证访问服务时除 Access Key ID 与 Secret Access Key 外还需提供 Session TokenEndpoint URLAWS 服务的自定义端点 URL默认端点由 SDK/CLI 按区域自动使用仅特殊场景需要覆盖Profile Name命名配置文件named profile如需使用default之外的配置文件时填写Assume Role ARN跨账户访问时扮演的角色 ARN使用AssumeRole时必填要求账户管理员为调用者附加允许AssumeRole的策略Assume Role Session Name扮演角色会话的标识符默认值为OpenMetadataSessionAssume Role Source Identity调用AssumeRole操作的主体指定的源身份可在 AWS CloudTrail 日志中用于审计是谁以某角色执行了操作。Azure Auth ConfigClient IDService Account 的唯一标识可在 Service Account 密钥文件中client_id键对应的值获取Client Secret在 Azure 门户 →App registrations→ 选择应用 →Manage → Certificates secrets → Client secrets中创建并复制Value列的值Tenant ID在应用的Overview部分复制Directory (tenant) IDStorage Account Name存储账户名称Key Vault NameKey Vault 名称Scopes在Expose an API部分复制Application ID URI并确保 URI 以/.default结尾不是则手动追加Host PortMySQL 实例地址格式hostname:port如localhost:3306OpenMetadata ingestion 运行在 Docker 中而服务在 localhost 时使用host.docker.internal:3306Database NameOpenMetadata 的层级为Database Service Database Schema Table。MySQL 没有独立的 Database 概念若希望数据归入非default的数据库名可在此指定Database Schema可选参数。设置后将元数据读取限制在单个数据库留空则扫描所有数据库SSL CAcaCertificateSSL 校验用 CA 证书ssl_caSSL CertificatesslCertificate客户端认证用 SSL 证书ssl_certSSL KeysslKey与 SSL 证书关联的私钥ssl_keyConnection Options构建连接 URL 时可发送给服务的附加连接选项Connection Arguments连接时可发送给服务的附加连接参数如安全或协议配置。Postgres 连接第二处含完整扩展配置在 MSSQL 小节之后文档还给出了 Postgres 连接的完整扩展配置用于 Airflow 元数据库为 Postgres 的场景Username连接 Postgres 的用户名需具备读取 Postgres 全部元数据的权限Auth ConfigauthType同样支持Basic Auth密码、IAM 认证AWS 服务、Azure 认证Azure 服务三种类型IAM 与 Azure 的字段清单与上文 MSSQL 部分一致AWS Access Key ID、Secret Access Key、Region、Session Token、Endpoint URL、Profile Name、Assume Role ARN/Session Name/Source IdentityAzure 的 Client ID、Client Secret、Tenant ID、Storage Account Name、Key Vault Name、ScopesHost and PorthostPortPostgres 实例地址格式hostname:port如localhost:5432Docker 场景使用host.docker.internal:5432Databasedatabase初始连接的 Postgres 数据库。如需采集所有数据库将ingestAllDatabases设为 trueSSL ModesslMode如prefer、verify-ca、allow等。若使用 IAM 认证建议选择allow或根据用例选择其他选项SSL CAcaCertificateSSL 校验用 CA 证书sslrootcert。Postgres 只需 CA 证书Classification NameclassificationName默认情况下Postgres 策略标签在 OpenMetadata 中以PostgresPolicyTags分类可自定义分类名称采集后在 Tags 页面的 Classifications 列表中可见Ingest All DatabasesingestAllDatabases勾选后采集集群中所有数据库不勾选则只采集上述database指定的数据库中的表Connection Arguments连接时可发送给服务的附加连接参数Connection Options构建连接 URL 时可发送给服务的附加连接选项。SQLite 连接如果 Airflow 使用 SQLite 作为元数据库Username / Password连接 SQLite 的用户名与密码内存数据库in-memory留空Host PortSQLite 实例地址格式hostname:port如localhost:3306Docker 场景使用host.docker.internal:3306内存数据库留空Database可选参数。设置后限制元数据读取到单个数据库留空则扫描所有数据库Database ModeSQLite 数据库的运行模式默认:memory:Connection Options / Connection Arguments同前文说明。Backend Connection 与 Airflow 3.x 环境变量文档中提到Backend Connection仅用于从运行在 OpenMetadata 实例中的 DAG 提取元数据如 GCS Composer。其底层实现值得注意connection.py在Airflow 2.x下通过airflow.settings.Session获取现有引擎借用 Airflow 进程内的全局 engine不会被 dispose在Airflow 3.x下Airflow 阻止直接 ORM 访问因此改为基于环境变量DB_SCHEME、DB_USER、DB_PASSWORD、DB_HOST、DB_PORT、AIRFLOW_DB可选DB_PROPERTIES拼装 SQLAlchemy 连接 URL 并创建引擎缺失任一变量会抛出SourceConnectionException提示 Airflow 3.x 执行环境必须定义这些变量。测试连接与错误诊断源码级在 OpenMetadata UI 中配置完连接后点击Test底层会执行 connection.py 中AirflowChecks定义的三步检查CheckAccessREST 通道调用test_get_version()严格校验返回状态与 JSON避免任意网页通过连通性测试在此之前还会对 host 做一次 TCP 探测以快速失败数据库通道执行SELECT 1PipelineDetailsAccessREST 通道列出 1 个 DAGGET /dags?limit1数据库通道查询serialized_dag表TaskDetailAccessREST 通道读取 1 个 DAG 的 tasks数据库通道读取serialized_dag的任务载荷兼容 Airflow 2.2.52.3.0 的不同列结构并处理COMPRESS_SERIALIZED_DAGS开启时空列的情况。同时AIRFLOW_ERRORS错误包提供了常见失败的诊断提示connection.pyHTTP 状态 / 异常诊断修复建议401认证失败检查用户名/密码或 Token是否正确且未过期403访问被拒绝凭证有效但权限不足为用户授予 Airflow REST API 读取权限404端点不存在Host and Port 指向 Airflow Web Server 而非 UI 或控制台页面429被限流稍后重试或提高 ingestion 账户的速率限制JSONDecodeErrorHost 不是 Airflow REST API返回的是网页而非 API检查地址并确认 REST API 已启用SSLErrorTLS 校验失败检查证书或在使用自签名证书时取消勾选 Verify SSLTimeout连接超时检查 Host and Port 及防火墙ConnectionError无法到达主机检查地址拼写与可达性数据库驱动错误access denied / password authentication failed且链路中存在 SQLAlchemyError数据库认证失败检查 Airflow 元数据库的用户名和密码小结OpenMetadata 的 Airflow Pipeline 连接器通过REST API 优先、数据库直连兜底的双通道设计覆盖了从自建单机到 Astronomer / Cloud Composer / MWAA 的全场景部署形态。配置时只需把握三条主线连通性Host and Port 必须指向 Airflow Web ServerREST 通道或元数据库DB 通道并注意 Docker 环境下的host.docker.internal主机名认证匹配按部署形态选择 Basic Auth / Access Token / GCP Service Account / MWAA 四选一AWS 与 Azure 凭证字段在 IAM / Azure 认证中复用血缘边界REST API 与数据库通道只提供 DAG 拓扑与运行状态表级、列级血缘必须另行部署 Lineage Backend 或 Lineage Operator。相关代码可在 ingestion/src/metadata/ingestion/source/pipeline/airflow/ 目录下深入阅读集成测试覆盖了全部认证方式与 DAG 提取流程见 test_airflow_api_connection.py单元测试则验证了连接构建与错误诊断逻辑见 test_airflow_connection.py。【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考