最近好几个朋友问我 OpenMetadata 连上 MySQL 之后,那些表结构、字段类型、主键索引到底是怎么被它自动摸清楚的。很多人配置完 service connection 之后,看着日志不知道自己机器里发生了什么,出了问题也没法下手。这篇文章我就拿 MySQL 举例,从代码角度把 OpenMetadata 获取数据库表元数据的完整链路拆开讲一讲,包括源码找法、关键类、底层 SQL、采集执行流程、常见坑和扩展思路。如果你正准备二次开发,或者排查元数据采集不准、连接失败之类的问题,应该会很有帮助。
1. 先搞懂OpenMetadata的元数据采集框架
1.1 为什么OpenMetadata选择插件化采集架构
OpenMetadata 本身的服务端是 Java 生态,但它的数据采集侧却选择了 Python 写的 Ingestion Framework。我第一次接触这个项目时也觉得奇怪,后来看多了才明白,这个设计是有意为之。元数据采集要对接的数据库太多了,MySQL、PostgreSQL、Oracle、Snowflake、BigQuery、Redshift,每种数据库的连接方式、方言、权限模型、元数据视图都不一样。如果用 Java 全部实现,光驱动和类型映射就能写出一大堆重复代码。
Python 生态里 SQLAlchemy 对数据库方言的适配做得非常成熟,几乎所有主流数据库都有对应的 dialect。OpenMetadata 只要基于 SQLAlchemy 写一套通用逻辑,再把每种数据库的特殊配置抽成独立的 Source 类,就能做到“核心通用、扩展简单”。所以你在源码里会看到大量source/database/目录下的文件,每个文件对应一种数据库,但真正的骨架逻辑都在公共基类里。
这个架构的另一个好处是采集进程可以脱离服务端单独跑。你可以把 ingestion 跑在本地,也可以接到调度系统里,完成后把结果通过 REST API 提交回 OpenMetadata 服务端。这样采集任务的压力不会压到控制面,不同环境的网络限制也好处理。理解这一点之后,再看代码就会发现,真正跟 MySQL 强相关的代码很少,大部分逻辑都在公共层。
1.2 表元数据究竟包含哪些信息,OpenMetadata又是怎么建模的
“表元数据”这四个字听起来很抽象,落到 MySQL 里其实就是information_schema里的那几张表。但很多初学者以为元数据就是CREATE TABLE语句里的字段定义,这是不够的。OpenMetadata 里一个 Table 实体除了表名、schema 名之外,还包含列信息、列类型、是否可空、默认值、注释、主键、外键、索引、表标签、所有者等等。
OpenMetadata 的元数据模型是分层的:Database -> DatabaseSchema -> Table -> Column。一个 MySQL 实例在 OpenMetadata 里通常被建模成一个 Database Service,下面可以包含多个 Database,再往下一层是 Schema。对 MySQL 来说,很多人会把database和schema混淆,因为 MySQL 里这两者基本是同一个概念。在 OpenMetadata 配置里有一个databaseSchema参数,指的就是你要采集的数据库名,最终会对应到DatabaseSchema实体。
列信息是采集的重点,OpenMetadata 的Column实体里包含name、dataType、dataLength、isNullable、constraint、defaultValue、description等字段。这些字段并不是 MySQL 直接给你一个 JSON,而是由采集程序把information_schema里的原始数据翻译成 OpenMetadata 的模型。后面我们讲代码时会看到,这个翻译过程是分步完成的,先由 SQLAlchemy 拿到标准结构,再由 OpenMetadata 转成自己的实体。
2. MySQL元数据获取的源码路径
2.1 源码仓库的目录结构,去哪里找MySQL相关代码
如果你是第一次看 OpenMetadata 源码,别直接搜SELECT * FROM,那样会很懵。它整个采集逻辑是分层的,建议按照下面的路径去读。我以开源仓库里ingestion/目录为基准:
| 关键文件路径 | 作用 |
|---|---|
ingestion/src/metadata/ingestion/source/database/mysql.py | MySQL 专属 Source 类,处理 MySQL 连接配置、类型映射 |
ingestion/src/metadata/ingestion/source/database/common_db_source.py | 关系型数据库公共逻辑,真正调用 SQLAlchemy inspector 的地方 |
ingestion/src/metadata/ingestion/source/database/sql_source.py | SQL 类数据源基类,定义采集流程框架 |
ingestion/src/metadata/ingestion/source/database/database_service.py | 数据库服务类,负责生成 DatabaseSchema、Table 实体 |
ingestion/src/metadata/ingestion/connections/database/mysql.py | 建立 MySQL 连接的工具类 |
ingestion/src/metadata/ingestion/connections/database/database_connection.py | 通用连接逻辑 |
ingestion/src/metadata/utils/sqlalchemy.py | SQLAlchemy 工具方法,创建 engine、inspector |
看代码时不要从mysql.py直接往下读,那样会陷入一堆细节。我建议先读common_db_source.py,因为这里的方法名非常直观,比如get_tables()、get_columns()、get_views(),一看就知道在干嘛。看完公共层再回头看mysql.py,你会发现它只重写了很小一部分,比如连接 URL 的拼法,还有少数类型转换的细节。
调用链大致是这样的:
MySQLSource -> CommonDbSource.get_tables() -> inspector.get_table_names() -> CommonDbSource.get_columns() -> inspector.get_columns() -> CommonDbSource.get_table_constraints() -> inspector.get_pk_constraint() -> inspector.get_foreign_keys() -> inspector.get_indexes()看到这里你应该已经有了感觉:真正跟 MySQL 底层的交互,全部封装在 SQLAlchemy 的 inspector 里,OpenMetadata 自己并不直接拼接 SQL。这种做法的好处是换数据库时,公共代码不用重写。
2.2 SQLAlchemy方言层:连接MySQL的桥梁
MySQL 在 OpenMetadata 里的连接逻辑集中在一个小文件connections/database/mysql.py中。它的核心任务是拼出一个合法的 SQLAlchemy URL,然后通过create_engine产生数据库连接。
常见的 URL 格式是:
mysql+pymysql://user:password@host:3306/database?charset=utf8mb4其中pymysql是 Python 连 MySQL 的驱动。OpenMetadata 还支持mysqlconnector,但实际项目里pymysql用得最多,因为安装简单、轻量。在代码中get_connection_url会根据 service connection 里的配置生成这个字符串。如果配置了databaseSchema,URL 里就会带上库名,但这并不意味着只采集这一个库,后续只是把这个库作为默认 schema 处理。
连接建立之后,OpenMetadata 会通过 SQLAlchemy 生成一个 inspector 对象:
from sqlalchemy import create_engine, inspect engine = create_engine(connection_url, connect_args={"connect_timeout": 10}) inspector = inspect(engine)inspector 是一个统一接口,对 MySQL 它会自动选择对应的 dialect,然后把get_columns这种高层调用翻译成 MySQL 能理解的information_schema查询。你不需要自己去解析 MySQL 的SHOW COLUMNS FROM table,SQLAlchemy 已经帮你完成了。
有一个细节值得注意:MySQL 8.0 默认认证插件是caching_sha2_password,如果用 pymysql 连接,偶尔会报cryptography包缺失。这个问题的本质是密码传输需要 RSA 加密,解决办法很简单,安装带加密依赖的版本,或者把 MySQL 用户改成mysql_native_password。这个问题我们在后面的坑位里细说。
2.3 抓取表元数据的底层SQL逻辑
虽然 OpenMetadata 自己很少直接写 SQL,但如果你想排查慢查询或者理解它做了什么,还是得知道它最终触发了哪些information_schema查询。以 MySQL 为例,SQLAlchemy inspector 在获取表信息时,底层大致会执行下面这些语句。
获取所有表名:
SELECT table_name FROM information_schema.tables WHERE table_schema = 'your_db';获取某张表的列信息,会查COLUMNS表:
SELECT column_name, column_type, is_nullable, column_default, character_set_name, comments FROM information_schema.COLUMNS WHERE table_schema = 'your_db' AND table_name = 'your_table' ORDER BY ordinal_position;获取主键约束:
SELECT column_name FROM information_schema.KEY_COLUMN_USAGE WHERE table_schema = 'your_db' AND table_name = 'your_table' AND constraint_name = 'PRIMARY';获取索引信息,会查STATISTICS表:
SELECT index_name, column_name, non_unique, seq_in_index FROM information_schema.STATISTICS WHERE table_schema = 'your_db' AND table_name = 'your_table' ORDER BY index_name, seq_in_index;OpenMetadata 拿到这些结果之后,会进行字段映射。例如,information_schema.COLUMNS里的column_type可能是int(11)、varchar(255)、enum(...)这种字符串,OpenMetadata 要把它解析成自己定义的数据类型。int会变成INT,varchar(255)会变成STRING且dataLength=255,datetime会变成DATETIME,json会变成JSON。这个过程在代码里对应CommonDbSource._get_columns和列类型映射函数。
所以你看,真正“抓元数据”的代码,很多精力不是花在 SQL 拼接上,而是花在结果解析和实体翻译上。理解了这一点,后面调试时你就知道该关注什么:如果采集出来的字段类型不对,多半是类型映射逻辑的问题;如果表清单不对,多半是过滤条件或者表名枚举问题。
3. 一次MySQL元数据采集的完整执行过程
3.1 最小权限准备与连接参数配置
在代码运行之前,连接能不能通、权限够不够是第一道门槛。OpenMetadata 采集 MySQL 元数据,最小权限其实很小,只需要能读目标库的表结构。理论上 MySQL 用户只要有访问目标库任意一张表的权限,information_schema里对应的元数据就能看到。但为了方便,很多人直接给一个只读账号:
CREATE USER 'openmetadata'@'%' IDENTIFIED BY 'strong_password'; GRANT SELECT ON *.* TO 'openmetadata'@'%'; FLUSH PRIVILEGES;GRANT SELECT ON *.*是全局只读,权限有点大,但很多内网环境图省事就这么干了。如果你只想让 OpenMetadata 采集某一个库,可以改成:
GRANT SELECT ON warehouse.* TO 'openmetadata'@'%';注意,如果后面要启用数据 Profiler,比如统计表的行数、空值率、数值分布,那需要的权限会更多,至少还要有在目标库上执行查询的权限。Profiler 会真的跑SELECT COUNT(*)之类的语句,所以只读权限可能还不够,要根据你的安全策略单独评估。
在 OpenMetadata 的配置文件中,连接参数一般长这样:
source: type: mysql serviceName: mysql_warehouse serviceConnection: config: type: Mysql username: openmetadata password: strong_password hostPort: localhost:3306 databaseSchema: warehouse sourceConfig: config: type: DatabaseMetadatadatabaseSchema字段非常关键,它指定了默认的数据库名。如果你有多个库要采集,通常做法是为每个库单独创建一个 service,或者通过 schema filter 分开处理。代码读到这个配置后,会用它来限定get_table_names的查询范围。
3.2 从pipeline配置到Workflow.run的关键代码走读
采集任务在 OpenMetadata 里的执行入口是Workflow。一个简单的工作流配置会被解析成一个 Python 字典,然后通过Workflow.create加载。核心代码结构类似:
from metadata.ingestion.api.workflow import Workflow config = { "source": { "type": "mysql", "serviceName": "mysql_warehouse", "serviceConnection": { "config": { "type": "Mysql", "username": "openmetadata", "password": "strong_password", "hostPort": "localhost:3306", "databaseSchema": "warehouse", } }, "sourceConfig": {"config": {"type": "DatabaseMetadata"}}, }, "sink": {"type": "metadata-rest", "config": {}}, "workflowConfig": { "openMetadataServerConfig": { "hostPort": "http://localhost:8585/api", "authProvider": "openmetadata", "securityConfig": { "jwtToken": "your_token", }, } }, } workflow = Workflow.create(config) workflow.execute() workflow.stop()execute内部是按“Source -> Processor -> Sink”的管道运行的。Source 负责从 MySQL 拉数据并转成 OpenMetadata 的实体,Processor 可以做一些扩展处理,比如过滤表、加标签,Sink 则负责把实体提交到 OpenMetadata 服务端。默认情况下,Processor 是透传的,不会改实体内容。
MySQLSource继承自CommonDbSource,next_record方法会不断返回Table实体。Table实体的构造过程在database_service.py里,它会先构建DatabaseSchema,再构建Table,再把Column、TableConstraint、Index等子对象挂到Table上。整个流程非常像流水线:源头是 MySQL 的 information_schema,终点是 OpenMetadata 的 API。
调试时你可以在workflow.execute()前后打日志,或者直接在配置里设置loggerLevel: DEBUG。这样运行metadata ingest -c your_config.yaml时,控制台会打印很多细节,包括连接信息、抓到的表数量、传给 Sink 的实体内容。看到实体内容后再去对照information_schema数据,就能很快定位是哪一层出了问题。
3.3 一条MySQL表记录在OpenMetadata里的字段流转
为了更直观地理解,我们拿一个常见的users表举例。假设 MySQL 里有这样一张表:
CREATE TABLE warehouse.users ( id INT NOT NULL AUTO_INCREMENT, email VARCHAR(255) COMMENT 'user email', is_active TINYINT(1) DEFAULT 1, created_at DATETIME, PRIMARY KEY (id), UNIQUE KEY uk_email (email) ) COMMENT='user table';OpenMetadata 采集时,information_schema会返回这张表的各类信息。经过 SQLAlchemy 之后,OpenMetadata 会把它们映射成下面这样的结构化实体:
| MySQL 信息 | SQLAlchemy/OpenMetadata 中间结果 | 最终字段 |
|---|---|---|
table_schema=warehouse | DatabaseSchema 名字 | DatabaseSchema.name |
table_name=users | Table 名字 | Table.name |
column_name=id | Column name | Column.name |
column_type=int | 类型解析为 INT | Column.dataType = INT |
is_nullable=NO | 非空 | Column.isNullable = False |
column_default/extra=auto_increment | 默认值和自增信息 | Column.defaultValue |
column_name=email | String 类型,长度 255 | Column.dataType = STRING,dataLength = 255 |
column_comment='user email' | 注释 | Column.description |
constraint_name=PRIMARY | 主键约束 | TableConstraint类型为 PRIMARY_KEY |
index_name=uk_email | 唯一索引 | Index实体,name=uk_email,columns=[email] |
这个映射过程不是一次性完成的。OpenMetadata 会先通过inspector.get_columns拿到列信息,再通过inspector.get_pk_constraint拿到主键,再通过inspector.get_foreign_keys和inspector.get_indexes拿到外键和索引。最后把这些数据组装成一个完整Table实体的列表。这也是为什么采集任务在表特别多的时候会比较慢,因为每个表都要执行好几条information_schema查询。
到这一步,你应该已经能回答“OpenMetadata 怎么获取数据库表元数据”这个问题了:它是通过 SQLAlchemy inspector 间接查 information_schema,然后做实体映射,最终通过 REST API 写入后端。理解了这条主线,后面的问题排查就会轻松很多。
4. 实际踩坑与排查经验
4.1 最常见的三类MySQL连接报错
连接层面的问题占了元数据采集故障的一大半。我总结几类典型报错和排查思路。
第一类:Access denied for user 'openmetadata'@'localhost'。这个报错除了密码错误之外,最常见原因是 MySQL 的用户 host 限制。你创建用户时如果用'openmetadata'@'localhost',但 OpenMetadata 服务从另一台机器连接,就会因为 host 不匹配被拒绝。解决办法是把用户改成'openmetadata'@'%',或者直接用能匹配来源 IP 的 host 创建用户。
第二类:Unknown database 'warehouse'。这个报错通常是配置里的databaseSchema写错了,或者 MySQL 端大小写敏感问题。Linux 下 MySQL 的库名是区分大小写的,如果实际库名是Warehouse而你配了warehouse,就会报这个错。可以用SHOW DATABASES;确认一下真实名称。
第三类:MySQL 8 默认认证插件导致的RuntimeError: 'cryptography' package is required for sha256_password or caching_sha2_password。这个问题很典型,因为 OpenMetadata 用 pymysql 连接 MySQL 8 时,默认认证插件是caching_sha2_password,需要cryptography包做 RSA 加密。解决办法有两种:一种是在 Python 环境里安装cryptography,另一种是把 MySQL 用户改成旧的认证方式:
ALTER USER 'openmetadata'@'%' IDENTIFIED WITH mysql_native_password BY 'strong_password';不过 MySQL 8.4 开始默认移除了mysql_native_password,所以更推荐直接装cryptography。在 OpenMetadata 的容器环境里,还需要确认 ingest 容器是否安装了对应依赖。
4.2 大库、多表采集慢的优化方案
如果你的 MySQL 实例里有几百甚至上千张表,全量采集时会明显感觉慢。原因不复杂:每张表都要查列、主键、外键、索引,还是串行执行的。而且information_schema的某些查询在表特别多、库特别多的时候,元数据本身也可能需要扫描,所以更要命。
我建议的优化方向有三个。第一,缩小采集范围。在 sourceConfig 里配置schemaFilterPattern和tableFilterPattern,把临时表、日志表、归档表排除掉。第二,拆流水线。不要把几百张表塞进一个 pipeline,可以按业务模块拆成多个 service 或者多个 ingestion,错峰执行,这样失败后定位也容易。第三,如果想要做增量采集,OpenMetadata 的元数据采集目前主要还是全量扫,但你可以通过只采集最近变更的表来降低压力,这个需要自己写点自定义逻辑。
另外要留意采集容器的内存。当表的列特别多时,OpenMetadata 要构造大量实体对象,内存会涨得很快。我遇到过一张表有一千多个字段的极端情况,直接 OOM。后来加了过滤规则,把这种异常表排除掉,或者单独设置 pipeline 才缓解。
4.3 类型映射不一致的处理办法
MySQL 的类型到 OpenMetadata 类型不是百分百一一对应。举个例子,MySQL 的enum类型在 SQLAlchemy 里可能返回ENUM,OpenMetadata 对它的处理有时候会落到STRING或者ENUM,这取决于版本。再比如tinyint(1),业务上经常用来表示布尔值,但底层返回的还是TINYINT,OpenMetadata 不会自动按布尔值处理。
如果发现类型映射不准,最直接的办法是继承 MySQLSource,重写类型转换方法。代码里通常有一个_get_column_type或者类似的方法,你可以针对特定类型做覆盖:
from metadata.generated.schema.entity.data.table import DataType from metadata.ingestion.source.database.mysql import MySQLSource class CustomMySQLSource(MySQLSource): def _get_column_type(self, column_info): col_type = column_info.get("type", "").lower() if col_type.startswith("tinyint(1)"): return DataType.BOOLEAN return super()._get_column_type(column_info)这种定制场景在银行、电商系统里非常常见,因为历史表设计五花八门。你在做二次开发时,不要想着改公共代码,而是尽量在自定义 Source 里覆盖方法,这样升级 OpenMetadata 版本时冲突最少。
4.4 验证采集结果与Debug日志
采集完怎么确认结果对不对?最直观的是去 OpenMetadata UI 搜索表名,看字段、约束、索引是否和源库一致。但 UI 只能看到最终结果,出问题时还得看日志。
推荐做法是先把workflowConfig里的loggerLevel调成DEBUG,然后跑一次metadata ingest -c your_config.yaml。DEBUG 日志会打出很多 SQLAlchemy 的 SQL 和执行耗时,你能直接看到它查了哪些information_schema表。如果连表清单都不对,就能马上发现是不是 schema 过滤写错了。
另一个技巧是开启 MySQL 的通用日志,确认 OpenMetadata 连接账号到底执行了什么语句。当然这个操作需要 DBA 配合,生产环境慎用。个人开发环境可以用performance_schema或直接在连接参数里加echo=True,SQLAlchemy 会把所有 SQL 打到控制台,省去很多猜测。
5. 将MySQL采集能力扩展到你自己的场景
5.1 从代码角度理解Source/Processor/Sink的扩展点
OpenMetadata 的采集代码最值得借鉴的地方,就是它是围绕 Source、Processor、Sink 三个接口设计的。Source 负责产出实体,Processor 负责加工实体,Sink 负责输出实体。要做定制,优先考虑 Processor,而不是去大改 Source。
举个例子,如果你希望采集表的时候自动给某些敏感字段打标签,那么可以写一个自定义 Processor,在拿到Table实体之后,遍历 columns,根据列名正则命中手机号、身份证号等,然后附加上Tag,再传给 Sink。源码里 Processor 的接口很简单,核心是一个process(record)方法,输入是 Source 产出的实体,输出是加工后的实体。
如果你要改元数据抓取逻辑,比如过滤某些特殊表、合并多个库的表,那么可以继承 MySQLSource。但注意CommonDbSource的方法设计得比较细,尽量只覆写单点方法,比如get_tables或get_columns,不要整个next_record重写,否则维护成本很高。我在一个项目里为了排除 MySQL 视图里的临时表,只重写了get_views方法,改动很小,效果却很明显。
5.2 从MySQL扩展到其他数据库需要动哪些地方
如果你熟悉了 MySQL 的采集链路,扩展到 PostgreSQL 其实是顺理成章的。OpenMetadata 已经内置了很多数据库 Source,比如postgres.py、trino.py、clickhouse.py等。新接入一个数据库时,主要工作是三件事。
第一,确认 SQLAlchemy 有没有对应的 dialect。比如 PostgreSQL 是postgresql+psycopg2://,ClickHouse 是clickhouse+http://。如果 SQLAlchemy 不支持,OpenMetadata 就要自己写一层适配,工作量会大很多。第二,实现get_connection_url和连接配置类,处理 SSL、连接池等差异。第三,根据目标数据库的元数据视图调整查询方式。PostgreSQL 查的是information_schema.columns,但某些字段的语义和 MySQL 不完全一致,需要在类型映射里补一些细节。
大多数情况下,基础表、列、索引都能通过 SQLAlchemy 通用接口拿到。真正需要额外处理的是那些数据库独有的特性,比如 MySQL 的auto_increment、PostgreSQL 的序列、ClickHouse 的引擎信息。OpenMetadata 在处理这些差异时,通常会在对应的 Source 类里加额外方法,而不是强行改公共模型。这也是一个很好的设计参考:公共模型保持通用,特殊细节在各数据源内部消化。
如果你看完这篇文章准备动手看源码,我建议你按这个顺序来:先跑通一个最简单的 MySQL 元数据采集,然后打开 DEBUG 日志,看 SQLAlchemy 实际执行了什么,接着到common_db_source.py里打断点,跟踪一个 Table 对象是怎么被造出来的。走完一遍之后,你对整个元数据采集的理解会明显上一个台阶,后面遇到其他数据源也能快速上手。