- 数据目录
- 数据治理
- 数据血缘
- 后端
- 前端
- 数据工程
- 数据集成
【免费下载链接】datahub
The Context Platform for your Data and AI Stack
导读
sqlalchemy是 DataHub 元数据接入框架中面向"一切具备 SQLAlchemy 方言的数据库"的通用型接入源。当目标数据库没有专属连接器、但社区已为其实现 SQLAlchemy 方言时,你可以通过本模块直接完成数据库/表/视图、字段、容器、血缘、数据画像与有状态删除检测的全套元数据摄入。读完本文,你将掌握该接入源的适用场景、配置全参、底层实现原理与排障路径,并能在几分钟内编写出可运行的 ingestion recipe。
一、模块定位:何时使用 sqlalchemy 接入源
根据仓库中该模块的说明文档(metadata-ingestion/docs/sources/sqlalchemy/sqlalchemy_pre.md):
The
sqlalchemymodule ingests metadata from SQLAlchemy into DataHub. It is intended for production ingestion workflows.
该接入源专为生产级摄入工作流设计,最典型的适用场景是:仓库中不存在针对你所用数据库的预构建 source,但该数据库已有第三方实现的 SQLAlchemy 方言(dialect)。此时只需自行pip install对应的方言包,即可复用本模块完成元数据摄入。
从源码看,接入源本体定义在 sql_generic.py:SQLAlchemyGenericSource继承自SQLAlchemySource,其注释明确写道:
- 使用 SQLAlchemy reflection(反射)机制发现 schema 元数据;
- 需要用户自行安装对应的方言包;
- 平台名通过配置项
platform指定; - 完整继承
SQLAlchemySource的全部能力(数据画像、过滤、域)。
模块标注为@support_status(SupportStatus.GA),即正式发布(GA)状态。
二、模块能提取哪些元数据
根据官方文档与本仓库实现,sqlalchemy接入源覆盖的核心元数据实体包括:
- 数据集(Dataset):数据库中的表(table)与视图(view);
- SchemaField:每张表的列及其类型、可空性、注释等字段元数据;
- 容器(Container):数据库(database)与模式(schema)层次结构;
- 表级与列级血缘(Lineage):视图到视图、表到视图的血缘关系;
- 数据画像(Profiling):可选的表、行、列统计信息;
- 有状态删除检测(Stateful Deletion Detection):通过 stateful ingestion 识别已消失的实体。
sqlalchemy_pre.md中归纳为三条核心提取项:
- 数据库、模式、视图和表的元数据;
- 每张表关联的列类型;
- 通过可选的 SQL profiling 提供的表、行、列统计信息。
三、底层工作原理:从 SQLAlchemy 反射到 DataHub 实体
3.1 核心类层次
接入源实现位于metadata-ingestion/src/datahub/ingestion/source/sql/目录,关键文件如下:
| 文件 | 职责 |
|---|---|
| sql_generic.py | 定义SQLAlchemyGenericConfig与SQLAlchemyGenericSource(平台名固定为sqlalchemy) |
| sql_common.py | 定义SQLAlchemySource,封装了 reflection、schema/字段/容器/血缘/画像的完整提取逻辑 |
| sql_config.py | 定义SQLCommonConfig、SQLFilterConfig、SQLAlchemyConnectionConfig等全部配置模型 |
| sqlalchemy_uri.py | 提供make_sqlalchemy_uri/parse_host_port工具函数,负责拼接连接 URI |
类继承关系可概括为:
SQLAlchemyGenericSource → SQLAlchemySource → StatefulIngestionSourceBase SQLAlchemyGenericConfig → SQLCommonConfig → StatefulIngestionConfigBase + PlatformInstanceConfigMixin + EnvConfigMixin + ...3.2 Reflection 发现机制
SQLAlchemySource基于 SQLAlchemy 官方的create_engine、inspect、Inspector反射 API 工作(见 sql_common.py):
from sqlalchemy import create_engine, inspect, log as sqlalchemy_log from sqlalchemy.engine.reflection import Inspector通过反射拿到数据库/模式/表/视图清单与列定义后,再借助 mce_builder 中的make_data_platform_urn、make_dataset_urn_with_platform_instance、make_schema_field_urn等函数,将原始对象转换为 DataHub 的 URN 与 MCP(MetadataChangeProposal),从而生成可被 GMS 消费的工作单元(WorkUnit)。
3.3 连接 URI 的生成
配置模型SQLAlchemyConnectionConfig.get_sql_alchemy_url()(sql_config.py)实现了两套建连方式:
- 直接传入
sqlalchemy_uri(优先级更高); - 或通过
scheme + username + password + host_port + database组合,调用make_sqlalchemy_uri()自动拼接。
其中parse_host_port(sqlalchemy_uri.py)负责解析host:port,支持端口缺失时回退默认端口、端口非法时静默使用默认值等边界处理。通用型SQLAlchemyGenericConfig则直接要求必填connect_uri,并把platform作为 URN 构造中的平台名。
四、快速开始:第一个 sqlalchemy recipe
仓库为模块提供了最小可运行配方模板 sqlalchemy_recipe.yml:
source: type: sqlalchemy config: # Coordinates connect_uri: "dialect+driver://username:password@host:port/database" sink: # sink configs实际运行时,type必须为sqlalchemy(对应@platform_name("SQLAlchemy", id="sqlalchemy")声明,见 sql_generic.py)。
以 PostgreSQL 为例(需先pip install psycopg2-binary或psycopg):
source: type: sqlalchemy config: platform: "postgres" connect_uri: "postgresql+psycopg2://datahub:datahub@localhost:5432/datahub" # 可选:过滤、画像、血缘等高级配置(见下文) sink: type: "datahub-rest" config: server: "http://datahub-gms:8080"随后执行:
datahub ingest -c sqlalchemy_recipe.yml注意:
platform配置项决定了 URN 中的平台名,应与实际数据库类型保持一致,便于 DataHub 侧平台识别与血缘关联。
五、配置详解:全部可调参数
SQLAlchemyGenericConfig继承自SQLCommonConfig(sql_config.py),完整配置项如下。
5.1 连接与坐标
| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
connect_uri | str | 是 | 连接 URI,格式见 SQLAlchemy 官方 database-urls 说明 |
platform | str | 是 | 摄入的平台名,用于构造 URN,例如postgres、mysql、clickhouse |
options | dict | 否 | 透传给SQLAlchemy.create_engine的 kwargs;如需设置 URL 中的连接参数,放在connect_args下 |
include_views | bool | 否(默认true) | 是否摄入视图 |
include_tables | bool | 否(默认true) | 是否摄入表 |
include_table_location_lineage | bool | 否(默认true) | 若源支持,摄入表到底层存储位置的血缘 |
include_view_lineage | bool | 否(默认true) | 使用 DataHub 的 SQL parser 填充视图→视图、表→视图血缘 |
include_view_column_lineage | bool | 否(默认true) | 基于 SQL parser 填充视图→视图、表→视图的列级血缘;依赖include_view_lineage开启 |
use_file_backed_cache | bool | 否(默认true) | 是否使用文件后备缓存存储视图定义 |
5.2 过滤模式(schema / table / view)
SQLFilterConfig(sql_config.py)提供三层正则过滤:
schema_pattern:按模式名(schema)过滤,例如analytics匹配 analytics 模式下所有表;在 schema 数量巨大时,用它可以避免无谓地拉取表清单后再过滤,是一种性能优化手段;table_pattern:按database.schema.table全名过滤,例如Customer.public.customer.*匹配 Customer 库 public schema 下所有 customer 开头的表;view_pattern:按database.schema.view全名过滤;未显式指定时默认继承table_pattern(源码中的view_pattern_is_table_pattern_unless_specified模型校验器实现了该行为)。
三者均使用AllowDenyPattern(allow 与 deny 正则列表)。另有:
profile_pattern:指定参与数据画像的表/列正则过滤,注意只有通过table_pattern的表才会被纳入画像;domain:Dict[str, AllowDenyPattern],按正则把数据库/模式/表挂到业务域(domain key 可以是 URN 如urn:li:domain:ec428203-ce86-4db3-985d-5a8ee6df32ba,也可以是 "Marketing" 这样的名称,DataHub 会自动解析为 URN,解析失败会报错)。
示例:
source: type: sqlalchemy config: platform: "postgres" connect_uri: "..." schema_pattern: allow: - "public" - "analytics" deny: - "information_schema" - "pg_.*" table_pattern: allow: - ".*\\.public\\.customer.*" profile_pattern: allow: - ".*\\.public\\.customer.*"5.3 数据画像(Profiling)
profiling字段类型为GEProfilingConfig(基于 Great Expectations 体系,见 ge_profiling_config.py)。启用画像需要同时满足:profiling.enabled = true,且operation_config判定画像开启(is_profiling_enabled()方法,sql_config.py)。
源码中ensure_profiling_pattern_is_passed_to_profiling模型校验器会将profile_pattern自动注入 profiling 配置,确保过滤一致性。画像覆盖表/行/列统计信息,如行数、空值率、去重值数等。
profiling: enabled: true operation_config: lower_limit: 10000 upper_limit: 10000005.4 状态化摄入(Stateful Ingestion)与删除检测
stateful_ingestion字段类型为StatefulStaleMetadataRemovalConfig(sql_config.py),继承自StatefulIngestionConfigBase。启用后,接入源会记录上次摄入的状态,并在下次运行时通过StatefulStaleMetadataRemovalHandler检测源端已不存在的表/视图等实体,生成删除类 MCP,实现有状态删除检测。
stateful_ingestion: enabled: true remove_stale_metadata: true5.5 其他继承能力
platform_instance:通过PlatformInstanceConfigMixin支持平台实例;env:通过EnvConfigMixin指定环境(PROD / DEV 等);- 增量血缘:通过
IncrementalLineageConfigMixin支持; - 数据分类:通过
ClassificationSourceConfigMixin支持基于样本的数据分类。
六、概念映射:源概念到 DataHub 实体
关联文档 README.md 给出了源概念到 DataHub 概念的映射表(注:特定于 sqlalchemy 的映射细节仍待完善,下表为 DataHub 通用概念映射):
| Source Concept | DataHub Concept | Notes |
|---|---|---|
| Platform/account/project scope | Platform Instance, Container | 在平台上下文内组织资产 |
| Core technical asset (例如 table/view/topic/file) | Dataset | 主要摄入的技术资产 |
| Schema fields / columns | SchemaField | 在支持 schema 提取时包含 |
| Ownership and collaboration principals | CorpUser, CorpGroup | 由支持所有权与身份元数据的模块发出 |
| Dependencies and processing relationships | Lineage edges | 在支持并启用血缘提取时可用 |
结合源码可以进一步对应:数据库与模式会生成DatabaseKey/SchemaKey容器(见 sql_common.py 与sql_utils.py中的gen_database_container/gen_schema_container),表/视图映射为DatasetSnapshot,列映射为SchemaFieldClass与SchemaMetadataClass,血缘映射为UpstreamLineageClass/FineGrainedLineageClass。
七、能力、限制与排障
7.1 能力矩阵
根据 sqlalchemy_post.md 与源码装饰器声明,模块能力如下:
- Domain 支持:通过
domain配置项支持(@capability(SourceCapability.DOMAINS, ...)); - 数据画像:通过配置可选启用(
@capability(SourceCapability.DATA_PROFILING, ...)); - 表/视图/字段元数据、容器层次、视图血缘、列级血缘、状态化删除检测均为 GA 能力。
能力是否生效受源平台 API、权限与暴露的元数据约束,请以"Important Capabilities"能力表为准,个别能力可能需要额外配置。
7.2 限制
- 接入源依赖 SQLAlchemy 方言包的反射能力:方言暴露多少元数据,DataHub 就能摄入多少;方言未实现的反射接口(如视图定义、存储过程)将导致相应元数据缺失;
- 需要自行
pip install对应方言包并保证版本兼容; - 血缘质量取决于视图 SQL 的可解析性(走 DataHub sql parser);
- 画像能力受数据库权限(如表扫描权限)限制。
7.3 排障建议
关联文档给出的排障顺序为:先校验凭据、权限、网络连通性与范围过滤器(schema/table/view pattern),再检查摄入日志中的源特定错误并相应调整配置。常见问题对应如下:
- 连接失败:优先检查
connect_uri中的dialect+driver是否正确、方言包是否安装、host:port是否可达; - 摄入实体为空:检查
schema_pattern/table_pattern是否误过滤,include_tables/include_views是否被关闭; - 血缘缺失:确认
include_view_lineage开启;列级血缘还需include_view_column_lineage开启; - 画像未执行:确认
profiling.enabled为 true,且表通过了profile_pattern与table_pattern; - 删除检测未生效:确认
stateful_ingestion.enabled与remove_stale_metadata均已开启。
八、总结
sqlalchemy接入源是 DataHub 覆盖长尾数据库的"万能钥匙":只要目标系统存在 SQLAlchemy 方言,即可通过一个connect_uri接入数据库/表/视图、字段、容器、血缘、画像与删除检测的完整元数据链路。结合仓库内 sql_generic.py、sql_config.py、sql_common.py 的实现,你可以按需组合过滤、画像、血缘与状态化配置,将任意方言数据库快速纳入 DataHub 的统一元数据平台。
- 数据目录
- 数据治理
- 数据血缘
- 后端
- 前端
- 数据工程
- 数据集成
【免费下载链接】datahub
The Context Platform for your Data and AI Stack
相关推荐
DataHub 通用 SQLAlchemy 元数据摄取源(sqlalchemy source)实战指南
DataHub 通用 SQLAlchemy 元数据摄取源(sqlalchemy source)实战指南 导读 DataHub 内置了面向常见数据库(Snowfl
数据目录数据治理数据血缘后端前端数据工程数据集成DataHub Apache Druid 元数据接入指南:SQLAlchemy 连接、概念映射与血统提取实践
DataHub Apache Druid 元数据接入指南:SQLAlchemy 连接、概念映射与血统提取实践 本文以 DataHub 仓库中的 Apache D
数据目录数据治理数据血缘后端前端数据工程数据集成DataHub 元数据接入(Metadata Ingestion)完全指南:SDK、CLI 与 50+ 数据源连接器
DataHub 元数据接入(Metadata Ingestion)完全指南:SDK、CLI 与 50+ 数据源连接器 DataHub 的元数据接入框架(Meta
数据目录数据治理数据血缘后端前端数据工程数据集成
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考