news 2026/9/18 21:58:15

DataHub MariaDB 数据源接入指南:元数据、血缘与使用统计的完整实践

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
DataHub MariaDB 数据源接入指南:元数据、血缘与使用统计的完整实践

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:3306scheme固定为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_portstringlocalhost:3306MariaDB 主机与端口
usernamestring连接用户名
passwordstring连接密码
databasestring目标数据库;为空则枚举全部并受database_pattern过滤
auth_modeenumPASSWORDPASSWORD(标准账号密码)或AWS_IAM(AWS RDS IAM 认证)
aws_configobject默认AWS RDS IAM 认证配置,仅在auth_mode=AWS_IAM时使用
optionsobject透传给 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_patternAllowDenyPattern类型,按数据库名过滤,默认全部允许;源码中系统库information_schemaperformance_schemamysqlsys会被无条件排除(见 mysql.py 与_is_allowed_database)。
  • schema_patterntable_patternview_pattern:沿用 SQL 通用源的层级过滤模式。

存储过程

  • include_stored_procedures:boolean,默认true,是否采集存储过程;
  • procedure_patternAllowDenyPattern,按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.enabledfalse是否启用数据画像
profiling.profile_table_row_limitnull仅对估算行数小于该值的表做画像。行数来自information_schema.tables.table_rows(存储引擎统计,可能过期);设置为null表示不限
profiling.profile_table_size_limitnull仅对大小小于指定 GB 数的表做画像,大小取自information_schema.tables.data_lengthnull表示不限

这两项护栏值必须大于 0,否则配置校验会直接报错(ValueError: profile_table_row_limit must be greater than 0 (or null to disable filtering))。源码在画像候选生成时读取_table_rows_cachedataset_name_to_storage_bytes(由information_schema.tables全表扫描填充),超限的表会被记入profiling_skipped_row_limit/profiling_skipped_size_limit统计。

源码还会在运行结束后给出"画像耗时的表"建议:若某张表画像耗时超过 30 秒,会提示设置profile_table_row_limitprofile_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=ONlog_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);
  • 只解析SELECTINSERTUPDATEDELETEREPLACEWITHCALLMERGE等 DML 语句(见_DML_LEADING_KEYWORDS),SETSHOWCOMMIT等管理语句被跳过;
  • 用户名取自user_hostpriv_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=Truegenerate_query_usage_statistics=Truegenerate_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=ONlog_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/18 21:58:12

以太网交换机基础配置与故障排查四阶段指南

简介&#xff1a;本资源是一份面向网络初学者与IT运维人员的以太网交换机基础培训教材&#xff0c;系统梳理局域网核心设备的关键原理与实践要点。内容覆盖以太网标准&#xff08;IEEE 802.3&#xff09;、MAC地址机制、以太网帧格式&#xff08;含Ethernet II、802.2 LLC、802…

作者头像 李华
网站建设 2026/9/18 21:57:43

PyTorch实现Unet图像分割:网络结构详解与torchsummary可视化

简介&#xff1a;面向图像分割初学者与PyTorch入门者的一份PDF说明文档&#xff0c;重点讲解用PyTorch搭建Unet卷积神经网络&#xff0c;并结合torchsummary可视化模型结构。Unet常用于医学图像等场景的像素级分割任务&#xff0c;整体呈对称U形&#xff0c;左侧四个下采样阶段…

作者头像 李华
网站建设 2026/9/18 21:55:58

CloddsBot市场索引API:Kalshi/Manifold实时行情语义检索实战指南

CloddsBot市场索引API&#xff1a;Kalshi/Manifold实时行情语义检索实战指南 【免费下载链接】CloddsBot Open Source AI trading agent that operates autonomously across 1000 markets - Polymarket, Kalshi, Binance, Hyperliquid, Solana DEXs, 5 EVM chains. Scans for e…

作者头像 李华
网站建设 2026/9/18 21:52:38

设计自己的小传输协议:导论与概念详解

一、引言&#xff1a;为什么需要理解传输协议在大多数应用开发者的日常工作中&#xff0c;网络通信往往被抽象成一行简单的调用&#xff1a;打开一个 Socket&#xff0c;写入一些字节&#xff0c;再读取一些字节。对于使用 HTTP、gRPC、WebSocket 这类成熟协议的业务系统来说&a…

作者头像 李华