DataHub MySQL 元数据接入指南:从连接配置、Usage 统计到 Profiling 的完整实战
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
导读
本文以 DataHub 官方仓库中的 MySQL Source 文档(metadata-ingestion/docs/sources/mysql/README.md)为核心骨架,系统讲解如何通过 DataHub 元数据接入框架将 MySQL 实例中的表、视图、Schema 字段、容器(Container)、血缘(Lineage)、数据画像(Profiling)与使用统计(Usage Statistics)纳入 DataHub 元数据体系。你将掌握连接配置、权限授予、多库过滤、性能调优、AWS RDS IAM 认证、双通道 Usage 采集,以及 MySQL 特有的 Profiling 防护配置,并可结合实际源码(mysql.py)理解每个参数背后的实现原理。
模块定位与能力总览
MySQL Source(模块标识mysql)是 DataHub 元数据接入框架(metadata-ingestion)中面向生产环境的 SQL 源之一,在源码中以MySQLSource类实现,标注为GA(General Availability)支持状态。从装饰器声明(mysql.py)可以确认其完整能力矩阵:
| 能力项 | 支持情况 | 说明 |
|---|---|---|
| 数据集(表/视图) | ✅ 默认启用 | 基于 SQLAlchemy 反射,两级命名空间(database.table) |
| Schema 字段(列) | ✅ 默认启用 | 含类型映射,注册了 GEOMETRY/POINT 等空间类型 |
| 容器(Container) | ✅ 默认启用 | database 级别容器(Two-tier 模型) |
| Platform Instance / Domains | ✅ 默认启用 | platform_instance与domain配置项 |
| 数据画像(Profiling) | ⚙️ 可选 | profiling.enabled开启,另有 MySQL 专属防护参数 |
| Usage 统计 | ⚙️ 可选 | include_usage_statistics开启,支持两种历史来源 |
| 血缘(Lineage) | ✅/⚙️ | 视图血缘默认开启;查询级表/列血缘随 Usage 开启 |
| 存储过程 | ✅ 默认启用 | include_stored_procedures,默认true |
| 有状态删除检测 | ✅ 默认启用 | Stateful Ingestion 的 stale-entity 软删除 |
MySQL Source 继承自TwoTierSQLAlchemySource(two_tier_sql_source.py),采用「database → table」两级命名空间(无独立 schema 层),因此get_identifier返回schema.table形式(mysql.py),容器层级也相应简化为 database 一级。这与 PostgreSQL 等三级命名空间源存在差异,配置时需注意。
概念映射
原文档给出了 DataHub 通用概念映射表(具体到 MySQL 的映射细节以实际实现为准):
| 源端概念 | DataHub 概念 | 说明 |
|---|---|---|
| Platform/account/project scope | Platform Instance、Container | 在平台上下文中组织资产 |
| 核心技术资产(表/视图/主题/文件) | Dataset | 主要被接入的技术资产 |
| Schema 字段/列 | SchemaField | 支持 schema 提取时包含 |
| 所有权与协作主体 | CorpUser、CorpGroup | 由支持所有权与身份元数据的模块产出 |
| 依赖与处理关系 | Lineage edges | 血缘提取支持并启用时可用 |
前置条件与权限授予
在运行接入之前,需要为接入用户授予最小权限。原文档明确要求两条 GRANT:
-- 元数据与 Profiling 所需(必需) GRANT SELECT ON DATABASE.* TO 'USERNAME'@'%'; -- 视图定义所需(必需) GRANT SHOW VIEW ON DATABASE.* TO 'USERNAME'@'%';第一项是元数据反射和 Profiling 的基础:SQLAlchemy Inspector 需要读取information_schema与表结构。第二项用于获取视图定义,是视图级血缘(include_view_lineage)的前提。
从源码看,系统库会被自动过滤。_SYSTEM_SCHEMAS(mysql.py)包含information_schema、performance_schema、mysql、sys,_is_allowed_database在两层过滤中都会排除它们(mysql.py)。
多数据库接入:database 与 database_pattern
MySQL 是两级命名空间,接入范围由database与database_pattern共同控制,两者的语义差异容易踩坑:
database不设置:接入用户可见的所有数据库。底层逻辑见 get_inspectors:当database为空时调用inspector.get_schema_names()枚举全部库,再逐个按database_pattern过滤。database_pattern:用正则筛选数据库,推荐用于多库选择性接入:
database_pattern: allow: - "^db_one$" - "^db_two$"database设为单个库:只接入该库。注意原文档特别强调:*不是database的合法值,通配符请使用database_pattern。
此外,两级架构下旧字段schema_pattern已废弃,配置解析会自动重命名为database_pattern(two_tier_sql_source.py),新配置应直接使用database_pattern。
表与视图过滤
接入粒度还可进一步用表/视图正则收敛:
table_pattern: allow: - "db_name\\.public\\.customer.*" deny: - "db_name\\.public\\.temp_.*" view_pattern: allow: [] # 默认回落到 table_patternview_pattern的默认值会回落到table_pattern(sql_config.py);include_tables与include_views默认均为true,可分别开关。
基础接入配置:完整 Recipe 详解
原文档配套的 mysql_recipe.yml 是一份可直接运行的接入配方。以下逐段展开并补充源码依据:
source: type: mysql config: # 连接坐标 host_port: localhost:3306 database: dbname # 凭据 username: root password: example # 可选:从查询历史推导 usage 统计与查询级血缘。 # include_usage_statistics: true # usage_source: performance_schema # 默认;归一化 digest,无需额外设置 # usage_source: general_log # 字面 SQL + 按用户归属;需 general_log=ON, log_output=TABLE # 如 MySQL 需要 SSL: # options: # connect_args: # ssl_ca: "path_to/server-ca.pem" # ssl_cert: "path_to/client-cert.pem" # ssl_key: "path_to/client-key.pem" # AWS RDS IAM 认证(替代密码) # auth_mode: "AWS_IAM" # aws_config: # aws_region: us-west-2 # AWS 高级配置(profile、角色扮演、重试): # auth_mode: "AWS_IAM" # aws_config: # aws_region: us-west-2 # aws_profile: production # aws_role: "arn:aws:iam::123456789:role/DataHubRole" # aws_retry_num: 10 # aws_retry_mode: adaptive # auth_mode 为 "AWS_IAM" 时 password 字段被忽略 # AWS 凭据可通过 AWS CLI、环境变量或 IAM 角色配置 sink: # sink 配置(例如 datahub-rest / datahub-kafka)关键配置参数速查
基于 MySQLConnectionConfig 与 MySQLConfig 的字段定义,汇总如下:
| 参数 | 默认值 | 说明 |
|---|---|---|
host_port | localhost:3306 | MySQL 主机与端口 |
database | None | 指定单个库;不设置则接入可见全部库(配合database_pattern) |
database_pattern | allow_all | 数据库正则过滤(两级模型下替代 schema_pattern) |
username/password | 无 | 连接凭据;密码为 SecretStr 类型 |
scheme | mysql+pymysql | SQLAlchemy 方言,内部固定 |
auth_mode | PASSWORD | PASSWORD或AWS_IAM(枚举见 MySQLAuthMode) |
aws_config | 默认 boto3 凭据链 | 仅AWS_IAM时生效 |
options | {} | 透传给 SQLAlchemy 的 engine 参数,如连接池、SSLconnect_args |
include_stored_procedures | true | 是否接入存储过程 |
include_usage_statistics | false | 是否从查询历史生成 Usage 与查询级血缘 |
usage_source | performance_schema | 查询历史来源(见下文) |
email_domain | None | 仅general_log时,为非邮箱用户名追加域名 |
include_view_lineage | true | 基于视图定义的表→视图血缘 |
include_view_column_lineage | true | 视图级列血缘,依赖前者开启 |
profiling.enabled | false | Profiling 总开关 |
profile_pattern | allow_all | Profiling 的表/列过滤,继承table_pattern约束 |
连接串生成:未显式提供sqlalchemy_uri时,由get_sql_alchemy_url(two_tier_sql_source.py)基于scheme + username + password + host_port + database拼出 URI;也支持直接用sqlalchemy_uri覆盖。
连接并发与 Profiling 资源上限调优
这是 MySQL 源最容易在生产中踩到的坑。原文档明确指出:开启 Profiling 后,源会并发打开连接。一个数据库在被画像期间,最多可能占用pool_size(默认5)+max_overflow个连接,而max_overflow默认等于profiling.max_workers(其默认值为5 × CPU 核数,见 ge_profiling_config.py 中ThreadPoolExecutor默认线程数的说明)。
对设置了较低max_user_connections的副本实例,容易触发如下错误:
User 'USERNAME' has exceeded the 'max_user_connections' resource三种缓解手段
- 调低
profiling.max_workers:例如设为5,直接限制并发 Profiling 连接数; - 通过
options直接调节 SQLAlchemy 连接池:options: pool_size: 2 max_overflow: 5 - 关闭 Profiling:若只需 schema 与血缘元数据:
profiling: enabled: false
从源码实现看,get_inspectors中每个数据库使用独立的 SQLAlchemy engine,并在反射/画像完成后立即dispose()(mysql.py),其注释明确说明这是为了避免低max_user_connections限制下连接被长期占用的设计。
MySQL 专属 Profiling 防护参数
除通用 Profiling 配置外,MySQL 源还独有两个大表防护参数(MySQLProfilingConfig):
| 参数 | 默认值 | 说明 |
|---|---|---|
profile_table_row_limit | null(不限制) | 仅对预估行数小于该值的表做画像。预估来自information_schema.tables.table_rows,是存储引擎统计值,可能过期 |
profile_table_size_limit | null(不限制) | 仅对大小小于该值(GB)的表做画像。大小取information_schema.tables.data_length |
profiling: enabled: true profile_table_row_limit: 1000000 # 跳过超过 100 万行的表 profile_table_size_limit: 10 # 跳过超过 10 GB 的表实现细节:这两个参数通过add_profile_metadata的information_schema.tables全量扫描构建缓存,再由generate_profile_candidates在 Python 侧过滤候选表(mysql.py)。行为要点:
data_length同时驱动 Dataset 的sizeInBytes属性,两处口径一致;- 参数值必须大于 0,否则配置校验直接报错(
_validate_positive_limits);null表示不启用该过滤; - 若
information_schema查询失败(权限受限、代理改写查询等),会回退到只读data_length的三列查询:大小限制仍生效,但行数限制在该次运行中失效,并通过报告给出提示; - 若某 schema 下所有表都超限,报告会输出 "No tables passed the row/size guardrail" 提示,帮助运维定位画像被跳过的原因。
测试用例 test_mysql_profiling.py 覆盖了行/大小限制的主路径、回退路径及告警/提示分支,可作为理解该机制行为的参考。
Usage 统计与查询级血缘
开启include_usage_statistics: true后,源会从查询历史中推导使用统计与查询级表血缘。需要强调:查询级血缘只要开启 usage 就会产出,与include_view_lineage相互独立——后者只控制基于视图定义的(view-definition)血缘。
该行为在_create_aggregator(mysql.py)中有明确实现:启用 usage 时,SqlParsingAggregator会同时开启generate_lineage、generate_queries、generate_query_usage_statistics与generate_usage_statistics。
usage_source枚举(MySQLUsageSource)提供两种历史来源,各有取舍:
performance_schema(默认)
读取归一化 digest 表events_statements_summary_by_digest,对应 SQL 见 mysql.py。
- 要求:启用
statements_digestconsumer(performance_schema 默认启用),并授予GRANT SELECT ON performance_schema.* TO 'USERNAME'@'%'; - 特点:开销低、无需额外配置;但 Usage 计数是跨用户聚合的,且从上次重置(服务器重启或表被 TRUNCATE)起累计,因此首次开启 usage 后的第一次接入可能把历史期间的全部计数归到单个时间戳,表现为一次异常的大峰值(配置字段注释已明确警示);
- 实现:digest 行没有 actor,
ObservedQuery不带 user;COUNT_STAR作为usage_multiplier计入统计;系统库与database_pattern之外的库在读取时即被过滤(mysql.py)。
general_log
读取mysql.general_log中的字面 SQL 语句,带用户与时间戳,对应 SQL 见 mysql.py。
- 要求:
general_log=ON、log_output=TABLE,并授予GRANT SELECT ON mysql.general_log TO 'USERNAME'@'%'; - 特点:有 general log 本身的开销,但能提供按用户归属与精确查询文本;由于 general_log 没有独立的库名列,源码通过解析
Connect/Init DB/USE语句维护每个会话的当前库(LRU 上限 10,000 个会话,见 mysql.py),以此解析未限定库名的表引用; - 用户归属:
user_host字段解析出登录用户名;若登录名是 LDAP/数据库用户名而非邮箱,需设置email_domain(如corp.com)让 usage 映射到正确的 CorpUser(mysql.py);已形如邮箱的用户名则原样保留; - 语句过滤:只有 SELECT/INSERT/UPDATE/DELETE/REPLACE/WITH/CALL/MERGE 等 DML 语句才进入解析(mysql.py),SET/SHOW/COMMIT 等管理语句被跳过。
include_usage_statistics: true usage_source: general_log # 或 performance_schema(默认) email_domain: corp.com # general_log 模式下 LDAP 用户映射邮箱 usage: start_time: "2026-09-01T00:00:00" end_time: "2026-09-18T00:00:00"容错设计:查询历史读取失败(如 consumer 未启用、缺少授权)不会中断整个接入流程——元数据已经产出,仅记录 warning 后跳过 usage 阶段(mysql.py)。单元测试 test_mysql_usage.py 对 digest 行映射、系统库过滤、database_pattern 过滤、general_log 用户解析与 email_domain 追加等均有覆盖。
AWS RDS IAM 认证
MySQL 源支持 AWS RDS 的 IAM 认证,替代用户名/密码方式,实现免静态口令的接入。
前置准备
参考 AWS 官方 RDS IAM 数据库认证文档完成三步:
- 在 RDS 实例上启用 IAM database authentication;
- 创建使用 IAM 认证的数据库用户;
- 配置带
rds-db:connect权限的 IAM 策略。
配置方式
在 recipe 中设置auth_mode: "AWS_IAM",可选的aws_config用于指定凭据与区域(默认走 boto3 的标准凭据链):
auth_mode: "AWS_IAM" aws_config: aws_region: us-west-2更完整的 AWS 配置(profile、角色扮演、重试策略):
auth_mode: "AWS_IAM" aws_config: aws_region: us-west-2 aws_profile: production aws_role: "arn:aws:iam::123456789:role/DataHubRole" aws_retry_num: 10 aws_retry_mode: adaptive要点(均已在 recipe 注释与源码中明确):
auth_mode为AWS_IAM时,password字段被忽略;- AWS 凭据可通过 AWS CLI、环境变量或 IAM 角色提供;
host_port必须带端口:源码在初始化时解析host_port,端口缺失会直接抛ValueError;username同样为必填(mysql.py);- 实现上通过 SQLAlchemy
do_connect事件监听器,在每次建连时用RDSIAMTokenManager.get_token()注入临时令牌作为密码;由于 PyMySQL 要求启用 SSL,监听器会确保ssl参数被设置(mysql.py)。
单元测试 test_mysql_rds_iam.py 覆盖了配置解析、默认 aws_config、自定义端口与缺用户名报错等场景。
存储过程与血缘补充
MySQL 源默认接入存储过程(include_stored_procedures: true),可通过procedure_pattern正则过滤,匹配格式为database.schema.procedure_name(对两级模型即database.procedure_name)。
实现上从information_schema.ROUTINES读取ROUTINE_DEFINITION(mysql.py),并做了关键处理:原生 SQL 存储过程的EXTERNAL_LANGUAGE通常为 NULL,源码会将其回填为QueryLanguageClass.SQL,否则血缘提取器会静默跳过所有原生存储过程——这正是 MySQL/MariaDB 上存储过程血缘得以生效的底层保障。
常见问题排查
原文档给出的排查顺序,结合源码补充如下:
- 凭据与权限:确认接入用户具备
SELECT与SHOW VIEW权限;usage 开启时额外确认 performance_schema 或 general_log 相关权限; - 连通性:确认
host_port可达、SSL 参数(options.connect_args)正确; - 范围过滤:确认
database/database_pattern/table_pattern配置符合预期——注意database不接受*通配符; - 连接数超限:若出现
max_user_connections报错,按上文调低profiling.max_workers或通过options收缩连接池; - 查看接入日志:日志中 source-specific 错误会给出明确提示,例如 usage 读取失败时的 warning 会附上 "Ensure the statements_digest consumer is enabled..." 之类的修复提示(mysql.py)。
小结
MySQL 源是 DataHub 元数据接入中成熟度较高(GA)的两级命名空间 SQL 源:开箱即用地接入表、视图、列、容器与视图血缘;按需开启 Profiling(配合profile_table_row_limit/profile_table_size_limit防大表画像)与 Usage(performance_schema或general_log双通道);并支持 AWS RDS IAM 免密认证。掌握连接并发模型、database与database_pattern的边界语义,以及 usage 双通道的取舍,即可在生产环境稳定运行 MySQL 元数据接入流水线。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考