DataHub MariaDB 数据源接入指南:元数据、血缘与使用统计的完整实践
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
MariaDB 是广泛用于存储与查询分析型或操作型数据的数据库平台。DataHub 的 MariaDB 集成以 MySQL Source 为基础实现,用于抽取数据集/表/视图、Schema 字段、容器等核心元数据,并支持表级与列级血缘、数据画像以及有状态删除检测。本文以仓库中的 MariaDB Source 文档 为主体,结合 MariaDB 源码实现、MySQL 基类 与集成测试,完整讲解配置方式、能力矩阵、使用统计与血缘提取原理,帮助你用一份可复制的 recipe 将 MariaDB 接入 DataHub。
概览:DataHub 如何集成 MariaDB
MariaDB 是一个用于存储和查询分析型或操作型数据的数据平台,更多信息可参考官方 MariaDB 文档(mariadb.org)。DataHub 对 MariaDB 的集成覆盖了以下核心元数据实体:
- datasets / tables / views:将数据库中的表与视图作为 Dataset 实体入库;
- schema fields:抽取每张表/视图的字段与列级 Schema 信息;
- containers:按数据库/平台作用域组织资产,形成容器层级;
- lineage:支持表级与列级血缘;
- data profiling:可选的表级数据画像(行数、大小等统计);
- stateful deletion detection:有状态删除检测,用于发现已删除的实体。
从源码结构看,MariaDBSource 直接继承自MySQLSource,仅重写了get_platform()以返回平台标识mariadb,其配置类复用了MySQLConfig(见@config_class(MySQLConfig)装饰器)。也就是说,MariaDB 的绝大多数行为、配置项与底层查询逻辑都与 MySQL 源共享,这对熟悉 MySQL 接入的用户是显著优势。
概念映射:源概念到 DataHub 概念
虽然 MariaDB 专属的概念映射仍在完善中,但以下展示了 DataHub 的通用概念映射关系:
| 源概念 | DataHub 概念 | 说明 |
|---|---|---|
| Platform/account/project scope(平台/账号/项目作用域) | Platform Instance、Container | 在平台上下文内组织资产 |
| Core technical asset(如表/视图/主题/文件) | Dataset | 主要的被采集技术资产 |
| Schema fields / columns(字段/列) | SchemaField | 支持 Schema 抽取时包含 |
| Ownership and collaboration principals(所有者与协作者) | CorpUser、CorpGroup | 由支持所有权与身份元数据的模块发出 |
| Dependencies and processing relationships(依赖与处理关系) | Lineage edges(血缘边) | 当支持并启用血缘抽取时可用 |
对于 MariaDB 源而言,两张表(表/视图)会成为 Dataset,列会成为 SchemaField,数据库本身可作为 Container 组织资产。
快速开始:编写 MariaDB 采集 Recipe
仓库在 metadata-ingestion/docs/sources/mariadb/mariadb_recipe.yml 中提供了一份完整的示例 recipe,是实际接入的起点。完整内容如下:
source: type: mariadb config: # 连接坐标 host_port: localhost:3306 database: dbname # 凭据 username: root password: example # 可选:从查询历史中推导使用统计与查询级血缘 # include_usage_statistics: true # usage_source: performance_schema # 默认值;规范化摘要,无需额外配置 # usage_source: general_log # 字面 SQL + 按用户归因;需要 general_log=ON, log_output=TABLE # 如果需要用 SSL 连接 MariaDB: # options: # connect_args: # ssl_ca: "path_to/server-ca.pem" # ssl_cert: "path_to/client-cert.pem" # ssl_key: "path_to/client-key.pem" # sink configs(写入目标,例如 datahub-rest / file)要点说明:
type: mariadb是 DataHub Ingestion 框架中的源类型标识,对应 connector 注册表 中的mariadb条目;host_port默认值为localhost:3306,scheme固定为mysql+pymysql(见 MySQLConnectionConfig);database指定目标库;若留空则枚举全部数据库,并可用database_pattern过滤。
运行采集的命令为:
datahub ingest -c mariadb_recipe.yml集成测试中同样使用run_datahub_cmd(["ingest", "-c", f"{config_file}"])的方式执行采集,可参考 test_mariadb.py。
核心配置项详解
以下配置继承自MySQLConfig(定义于 mysql.py),对 MariaDB 源完全适用。
连接与认证
| 配置项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
host_port | string | localhost:3306 | MariaDB 主机与端口 |
username | string | 无 | 连接用户名 |
password | string | 无 | 连接密码 |
database | string | 无 | 目标数据库;为空则枚举全部并受database_pattern过滤 |
auth_mode | enum | PASSWORD | PASSWORD(标准账号密码)或AWS_IAM(AWS RDS IAM 认证) |
aws_config | object | 默认 | AWS RDS IAM 认证配置,仅在auth_mode=AWS_IAM时使用 |
options | object | 无 | 透传给 SQLAlchemy 引擎的额外选项,如connect_args中的 SSL 参数 |
关于auth_mode=AWS_IAM:源码在初始化时通过parse_host_port解析主机与端口(端口必填),并通过RDSIAMTokenManager生成临时令牌,在每次数据库连接时通过 SQLAlchemy 的do_connect事件监听器注入password;PyMySQL 要求 RDS IAM 认证必须启用 SSL(见 mysql.py 与 mysql.py)。
过滤模式
database_pattern:AllowDenyPattern类型,按数据库名过滤,默认全部允许;源码中系统库information_schema、performance_schema、mysql、sys会被无条件排除(见 mysql.py 与_is_allowed_database)。schema_pattern、table_pattern、view_pattern:沿用 SQL 通用源的层级过滤模式。
存储过程
include_stored_procedures:boolean,默认true,是否采集存储过程;procedure_pattern:AllowDenyPattern,按database.schema.procedure_name全名匹配过滤,例如Customer.public.customer.*匹配 Customer 库 public schema 下所有以 customer 开头的存储过程。
采集存储过程时,源码从information_schema.ROUTINES读取过程名、定义与语言;对于原生 SQL 存储过程,EXTERNAL_LANGUAGE为 NULL,会默认设为QueryLanguageClass.SQL,从而保证过程级血缘提取能够正常触发(见 mysql.py)。集成测试 mariadb_to_file.yml 演示了通过procedure_pattern只采集test_db.*并排除.*_temp$的过滤方式。
数据画像(Profiling)
profiling配置复用GEProfilingConfig,并额外提供两个 MariaDB/MySQL 特有的护栏参数(见 MySQLProfilingConfig):
| 配置项 | 默认值 | 说明 |
|---|---|---|
profiling.enabled | false | 是否启用数据画像 |
profiling.profile_table_row_limit | null | 仅对估算行数小于该值的表做画像。行数来自information_schema.tables.table_rows(存储引擎统计,可能过期);设置为null表示不限 |
profiling.profile_table_size_limit | null | 仅对大小小于指定 GB 数的表做画像,大小取自information_schema.tables.data_length;null表示不限 |
这两项护栏值必须大于 0,否则配置校验会直接报错(ValueError: profile_table_row_limit must be greater than 0 (or null to disable filtering))。源码在画像候选生成时读取_table_rows_cache与dataset_name_to_storage_bytes(由information_schema.tables全表扫描填充),超限的表会被记入profiling_skipped_row_limit/profiling_skipped_size_limit统计。
源码还会在运行结束后给出"画像耗时的表"建议:若某张表画像耗时超过 30 秒,会提示设置profile_table_row_limit、profile_table_size_limit或调低profiling.max_workers(并发全表扫描会倍增峰值内存,建议例如max_workers=5以缓解内存压力,见 mysql.py)。注意profile_if_updated_since_days对 MySQL/MariaDB 画像不生效,会输出提示后被忽略。
使用统计与查询血缘:两种查询历史来源
这是 MariaDB 集成最具特色的能力。开启include_usage_statistics: true后,DataHub 会读取查询历史,同时产出使用统计(datasetUsageStatistics)与查询级表血缘。查询历史的来源由usage_source控制,二选一:
方式一:performance_schema(默认,零配置)
source: type: mariadb config: host_port: localhost:3306 database: dbname username: root password: example include_usage_statistics: true usage_source: performance_schema # 默认值底层查询performance_schema.events_statements_summary_by_digest(见 mysql.py),每行是一个规范化语句摘要:
DIGEST_TEXT:去除了字面量(替换为?)的语句文本;COUNT_STAR:自上次计数器重置以来(服务器重启或表被 truncate)的执行次数;LAST_SEEN:最近一次执行时间;SCHEMA_NAME:语句所在库。
特点与注意点:
- 无需任何服务端配置,只要 statements digest 消费者已启用、用户对 performance_schema 有 SELECT 权限即可;
- 无按用户归因:摘要按语句聚合,不携带执行者,因此不会产生 per-user 使用统计(集成测试断言
userCounts为空,见 test_mariadb.py); - 首次采集可能产生峰值:由于
COUNT_STAR是累计值,启用后第一次采集可能把历史全部归到一个时间戳上,出现单日大峰值,属预期现象; - 每条摘要以
usage_multiplier=count计入聚合器。
方式二:general_log(字面 SQL + 按用户归因)
source: type: mariadb config: host_port: localhost:3306 database: dbname username: root password: example include_usage_statistics: true usage_source: general_log # 需要服务端开启 general log要求 MariaDB 服务端开启general_log=ON且log_output=TABLE,并且连接用户对mysql.general_log表有 SELECT 权限。底层查询见 mysql.py:
SELECT event_time, user_host, thread_id, command_type, CONVERT(argument USING utf8mb4) AS argument FROM mysql.general_log WHERE command_type IN ('Query', 'Init DB', 'Connect') AND event_time BETWEEN :start_time AND :end_time ORDER BY event_time, thread_id特点与注意点:
- 保留字面 SQL、执行用户与真实时间戳,因此可以按用户归因使用统计;
- 由于 general_log 没有 schema 列,源码通过跟踪每个会话(thread_id)的当前库来解析未限定表名:
Connect事件解析 "user@host on db using protocol" 中的初始库,Init DB/USE <db>语句切换当前库;会话映射采用 LRU 上限(默认最多跟踪 10000 个会话,见 mysql.py); - 只解析
SELECT、INSERT、UPDATE、DELETE、REPLACE、WITH、CALL、MERGE等 DML 语句(见_DML_LEADING_KEYWORDS),SET、SHOW、COMMIT等管理语句被跳过; - 用户名取自
user_host的priv_user[login_user] @ host [ip]格式;若用户名不是邮箱形式,可用email_domain配置项补全域名(如 LDAP 登录名),使其正确映射到 CorpUser,例如email_domain: example.com; Connect事件还用于记录会话初始默认库:当客户端连接时未选择数据库,<db>槽位为空,则该会话在出现USE/Init DB前 schema 未知,相关语句会被跳过(调试日志提示 "has no known database")。
两条来源路径均通过_usage_connection上下文管理器建立一个一次性、用完即释放(NullPool + dispose)的 UTC 时区连接(SET time_zone = '+00:00'),保证时间戳解析一致(见 mysql.py)。若读取查询历史失败(如消费者未启用、权限不足),仅记录 warning 并跳过使用统计,不会中断已完成的元数据采集。
血缘能力:视图血缘与查询血缘
根据 mariadb.py 的能力声明,MariaDB 源的血缘能力分为两档:
| 能力 | 默认状态 | 说明 |
|---|---|---|
| LINEAGE_COARSE(粗粒度) | 默认开启 | 视图表级血缘由include_view_lineage控制;开启使用统计后,还可从查询历史推导表级血缘 |
| LINEAGE_FINE(细粒度/列级) | 默认开启 | 视图列级血缘由include_view_column_lineage控制;开启使用统计后,还可从查询历史推导列级血缘 |
实现上,MySQLSource._create_aggregator在开启include_usage_statistics时会构造一个"查询 + 使用统计 + 血缘"全功能的SqlParsingAggregator(见 mysql.py):
generate_lineage=True:查询历史推导表血缘,与include_view_lineage相互独立(后者只控制视图定义血缘);generate_queries=True、generate_query_usage_statistics=True、generate_usage_statistics=True:同时产出查询与使用统计;generate_operations=False:两种来源均不产出操作(operation)aspect;_is_allowed_table/_is_temp_table回调确保未被采集的表(临时表、被过滤库的表、解析器误展开的db.db.table引用)仅作为血缘中转节点,不会产生"幽灵数据集"。
另外,源码会保存发现的 Schema 到解析器(即使未启用视图血缘,只要开启使用统计也需要,见_save_schema_to_resolver),保证未限定的表引用能被正确解析。
能力矩阵小结
以下能力来自 mariadb.py 的装饰器声明:
| 能力 | 支持情况 |
|---|---|
| Platform Instance | 默认启用 |
| Domains | 通过domain配置字段支持 |
| Data Profiling(数据画像) | 可选,通过配置启用 |
| Usage Stats(使用统计) | 可选,通过include_usage_statistics开启;读取 performance_schema 摘要(默认)或 mysql.general_log(usage_source: general_log),同时产出查询级表血缘 |
| Lineage - Coarse(粗粒度) | 视图默认启用(include_view_lineage);表级血缘可在开启使用统计后由查询历史推导 |
| Lineage - Fine(列级) | 视图默认启用(include_view_column_lineage);列级血缘可在开启使用统计后由查询历史推导 |
MariaDB 源在 mariadb.py 中标注的 SupportStatus 为GA(正式可用)。
测试验证:集成测试与单元测试
仓库为 MariaDB 源提供了完整的测试覆盖,是理解行为边界的绝佳参照:
- 单元测试test_mariadb_source.py 验证
MariaDBSource.platform == "mariadb",确认平台标识正确设置; - 集成测试test_mariadb.py 使用
mariadb:11.4Docker 镜像(见 docker-compose.yml)执行三类场景:- 元数据采集:以
mariadb_to_file.yml为配置将结果写入 MCE JSON,与 golden 文件mariadb_mces_golden.json对比校验(覆盖存储过程采集与过程级血缘); - performance_schema 使用统计:断言生成了
raw_customer_data的使用统计、无 per-user 归因、查询血缘存在(见 test_mariadb.py); - general_log 使用统计:断言按用户归因(
urn:li:corpuser:root出现在 userCounts 中)且查询血缘存在(见 test_mariadb.py); - 连接测试:对可达与不可达的 host_port 分别验证
test connection成功与失败。
- 元数据采集:以
这些测试中执行的使用工作负载由tests/test_helpers/mysql_usage_helpers.py提供,与 MySQL 源共享。
常见问题与排障
- 使用统计未生成:检查
include_usage_statistics是否开启;performance_schema 场景需确认 digest 消费者已启用且用户对performance_schema有 SELECT 权限;general_log 场景需确认general_log=ON、log_output=TABLE且用户对mysql.general_log有 SELECT 权限。读取失败仅告警不中断。 - general_log 场景下部分语句被跳过:会话在出现
USE/Init DB/Connect前 schema 未知的语句会被跳过;属于预期行为。 - 首次启用使用统计出现单日大峰值:performance_schema 的
COUNT_STAR为累计值,首次采集会把历史全部归到一个时间戳,属已知现象,可等待计数器重置后再次采集。 - 画像未覆盖大表:确认
profile_table_row_limit/profile_table_size_limit是否为null或误设为 0/负数(会直接校验失败);行数为information_schema估算值,可能过期,可执行ANALYZE TABLE刷新。 - 存储过程血缘缺失:确认
include_stored_procedures为 true 且procedure_pattern未过滤掉目标过程;原生 SQL 过程的语言默认按 SQL 处理以支持血缘提取。
延伸阅读
- MariaDB 源文档:本文的原始依据;
- MariaDB 示例 Recipe:可直接复制的配置模板;
- MariaDB 源码实现:能力声明与平台标识;
- MySQL 基类实现:全部配置项与查询历史/画像/存储过程逻辑;
- MariaDB 集成测试 与 单元测试:行为边界的可执行验证;
- 通用 SQL 源概念映射可参考 DataHub 文档中的 metadata-model 与概念说明(docs/what 目录)。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考