- 数据目录
- 数据治理
- 数据血缘
- 后端
- 前端
- 数据工程
- 数据集成
【免费下载链接】datahub
The Context Platform for your Data and AI Stack
本篇技术指南以 DataHub 开源仓库metadata-ingestion/tests/performance下的性能测试框架为主题,系统讲解其设计思想、公共数据生成/数据模型组件、各数据源(Snowflake、BigQuery、Databricks、Unity Catalog)与 SQL 解析(sqlglot、SQL Aggregator)的基准测试实现,以及 GraphQL 投影与内存泄漏两类专项性能测试。读者读完可掌握如何运行python -m tests.performance.<test_name>,理解WorkUnit吞吐、峰值内存、PerfTimer计时等核心度量手段,并学会将其应用到自己的 ingestion source 性能验证与回归防护中。
一、框架概览:为什么 ingestion source 需要专门的性能测试
数据接入(ingestion)是 DataHub 元数据平台最重负载的环节之一。仓库级 source(如 Snowflake、BigQuery、Databricks Unity Catalog)动辄需要拉取数万张表、数万个视图、数十万条查询日志,任何一次"顺序遍历改成了 O(n²)""缓存失效""SQL 解析器升级引发内存泄漏"都可能让一次全量接入从分钟级劣化为小时级,或让进程内存暴涨数 GB。
metadata-ingestion/tests/performance目录正是为回答这类问题而存在的专用性能测试模块。其 README 开宗明义:
This module provides a framework for performance testing our ingestion sources.
它的统一运行方式非常简洁,任何测试都可以通过 Python 模块方式直接执行:
python -m tests.performance.<test_name> # 例如: python -m tests.performance.snowflake.test_snowflake从目录结构看,该框架由三部分构成:
- 公共组件:data_model.py(数据模型)、data_generation.py(数据生成器)、helpers.py(度量工具);
- 仓库型 source 基准测试:
snowflake/、bigquery/、databricks/三个子目录; - 专项性能测试:sql/test_sql_formatter.py、sql_parsing/test_sql_aggregator.py、test_graphql_projection_perf.py、test_sqlglot_memory_leak.py。
二、性能度量的三大支柱:WorkUnit 计数、峰值内存与 PerfTimer
先看框架的度量基础设施,这是所有测试输出口径的统一基础。
2.1 WorkUnit 消费器workunit_sink
helpers.py 中的workunit_sink是每个基准测试都要用的"模拟下游消费者":它不真正写入任何元数据服务,而是以极低开销遍历 source 产出的全部MetadataWorkUnit,同时统计两件事:
- 产出 WorkUnit 总数(即元数据产出量);
- 进程峰值 RSS 内存(通过
psutil读取当前进程/proc/self/statm对应的memory_info().rss)。
def workunit_sink(workunits: Iterable[MetadataWorkUnit]) -> Tuple[int, int]: peak_memory_usage = psutil.Process(os.getpid()).memory_info().rss i: int = 0 for i, _wu in enumerate(workunits): if i % 10_000 == 0: peak_memory_usage = max( peak_memory_usage, psutil.Process(os.getpid()).memory_info().rss ) peak_memory_usage = max( peak_memory_usage, psutil.Process(os.getpid()).memory_info().rss ) return i, peak_memory_usage注意实现细节:并非每消费一个 WorkUnit 就采样一次内存(那样会拖慢被测路径本身),而是每 10,000 个采样一次,并在收尾时再补采一次,取最大值作为峰值。这保证了"测量对被测代码的影响尽量小"。
MetadataWorkUnit是 DataHub ingestion 管线的核心产物类型,定义于 datahub.ingestion.api.workunit,一条 WorkUnit 代表一个可独立落地的元数据变更(如一个 dataset 快照、一条 lineage 关系或一条 usage 统计)。
2.2 计时器PerfTimer
所有基准测试用with PerfTimer() as timer: ...包裹被测代码,再通过timer.elapsed_seconds(digits=2)输出耗时。其实现位于 datahub/utilities/perf_timer.py:基于time.perf_counter()的上下文管理器,还支持pause()/进入暂停状态继续计时,能够精确测量(默认保留 4 位小数秒)并叠加暂停前的活跃时间,非常适合需要隔离"数据生成阶段"与"真正被测阶段"的基准脚本。
2.3 人类可读输出与报告
各测试统一使用humanfriendly.format_size()将字节数格式化为可读单位,并打印source.get_report().as_string()(source 的内部统计报告)以及 report 中的关键指标,如 BigQuery 测试打印的usage_state_size(usage 去重状态占用磁盘大小)与num_usage_query_hash_collisions(查询哈希碰撞数)。
三、公共组件:数据模型与合成数据生成器
仓库型 source 的性能测试不能依赖真实生产环境(既不安全也不可重复),因此框架先用程序化生成的方式构造一份接近真实的元数据宇宙,再喂给被测 source。
3.1 数据模型data_model.py
data_model.py 定义了用于构建模拟仓库的最小领域模型:
Container:层级容器(如 database → schema),支持parent指针形成多层嵌套;Column/ColumnType:列与列类型枚举(INTEGER、FLOAT、STRING、BOOLEAN、DATETIME,基于StrEnum实现);Table/View:表与视图,Table.columns用OrderedDict保持列序,upstreams记录血缘上游;View额外持有definition(视图定义 SQL);name_components属性沿容器链展开得到[catalog, schema, table]形式的全限定名,这与 BigQuery/Databricks 的标识符模型完全一致;FieldAccess:一次查询中"哪个表上的哪一列被访问",用于 usage 场景;Query:一条审计日志级别的查询记录,字段包括text(SQL 文本)、type(StatementType,取值SELECT/INSERT/UPDATE/DELETE/CREATE/ALTER/DROP/CUSTOM/UNKNOWN)、actor(执行用户)、timestamp、fields_accessed与可选的object_modified(被修改的对象)。
3.2 分布与数据生成data_generation.py
data_generation.py 是框架最核心的公共工具,它并不生成"平均分布"的死板数据,而是让数据特征贴近真实仓库的幂律形态:
Distribution抽象 + 两种实现:NormalDistribution(mu, sigma):正态分布,适合列数、查询长度等"集中在均值附近"的指标;LomaxDistribution(scale, shape):重尾分布(等价于pareto(scale, shape) - scale),适合模拟真实环境中的"少数热表被大量查询、极少数超大上游"现象——源码注释明确给出了该分布在血缘上游数量上的百分位形态(75th=0、80th=1、95th=2、99th=4、99.99th=15)。- 二者均可通过
sample(floor=, ceiling=)施加上下界裁剪,防止生成越界值。
generate_data(...)入口:按参数生成整套仓库元数据,关键参数包括:num_containers:可以是整数(单层)或列表(多层容器,如[1, 100, 5000]表示 1 个 metastore/目录层、100 个 catalog/schema 层、5000 个 schema/子层);num_tables/num_views:表与视图数量;columns_per_table、parents_per_view、view_definition_length:各维度分布;time_range:usage 时间窗(默认 14 天)。 它返回SeedMetadata:containers(分层容器列表)、tables、views、start_time/end_time。每张表自动带上id列(ID_COLUMN = "id"),便于后续生成 join 血缘。
generate_lineage:为每张表按 Lomax 分布采样上游数量,并优先让"自身上游很多"的表更可能成为别人的上游(factor = 1 + len(tables) // 10的加权抽样),模拟真实 DAG 中"中心表被反复依赖"的结构。generate_queries(...):批量生成查询日志,参数包括num_selects(纯 SELECT 数)、num_operations(写操作数)、num_unique_queries(去重后的 SQL 文本数)、num_users、tables_per_select、columns_per_select等分布。查询文本来自faker生成的自然语言段落(模拟无 Schema 时的"脏"SQL),操作语句会从OPERATION_TYPES随机挑选并携带object_modified。模块还提供了if __name__ == "__main__": generate_data(10, 1000, 10)的快速自检入口。
值得注意的是,模块 docstring 明确说明这是"work in progress,按需逐步构建",并展望了未来两种更真实的方案:对生产 DataHub 实例的数据做匿名化去重,或引入 Faker 生成更拟人的数据。因此在使用时应将其理解为合成数据的脚手架而非最终形态。
四、仓库型 Source 基准测试实战
这一节逐一拆解三个真实基准测试的构造方式、被测规模与输出口径。
4.1 Snowflake:全管线 mock 压力测试
test_snowflake.py 是原 README 示例中的测试。核心思路:用unittest.mock完全替换 Snowflake 驱动连接,把驱动层换成纯内存的"假数据库",从而在无真实 Snowflake 账号的情况下跑通SnowflakeV2Source的完整get_workunits()流程。
关键构造如下:
with mock.patch("snowflake.connector.connect") as mock_connect: sf_connection = mock.MagicMock() sf_cursor = mock.MagicMock() mock_connect.return_value = sf_connection sf_connection.cursor.return_value = sf_cursor sf_cursor.execute.side_effect = functools.partial( default_query_results, num_tables=30000, num_views=10000, num_cols=30, num_ops=30, num_usages=500, )default_query_results来自 tests/integration/snowflake/common.py,是集成测试共享的模拟查询结果集工厂,这里把它放大到30,000 张表、10,000 个视图、每表 30 列、30 条操作、500 条 usage的量级;SnowflakeV2Config配置了include_technical_schema=False、include_table_lineage=True、include_usage_stats=True、include_operational_stats=True、format_sql_queries=True,即同时压测技术元数据、血缘、usage 与 SQL 格式化四条链路;- 计时与度量沿用前文的
PerfTimer+workunit_sink,最后打印source.get_report().as_string()与source.report.aspects。
运行方式与 README 示例完全一致:
cd metadata-ingestion python -m tests.performance.snowflake.test_snowflake4.2 BigQuery:端到端 usage 事件流基准
test_bigquery_usage.py 是最复杂的基准之一,它覆盖了"合成数据 → 模拟审计事件 → BigQuery usage extractor → WorkUnit"的完整链路:
- 种子数据:
generate_data(num_containers=2000, num_tables=20000, num_views=2000, time_range=timedelta(days=7))生成 2 万表、2 千视图、2 千容器; - 事件生成:100 个项目随机分配给表,随后
generate_queries(..., num_selects=240_000, num_operations=800_000, num_unique_queries=50_000, num_users=2000, query_length=NormalDistribution(2000, 500))产生104 万条查询,按时间排序后经 bigquery_events.py 的generate_events转成AuditEvent(QueryEvent+ReadEvent)流; - 被测对象:
BigQueryUsageExtractor(来自 datahub/ingestion/source/bigquery_v2/usage.py),配置usage.include_top_n_queries=True、top_n_queries=10、apply_view_usage_to_tables=True,并设置file_backed_cache_size=1000(文件后备缓存上限); - 分阶段度量:
BigQueryV2Report.new_stage(...)把"Seed Data Generation / Event Generation / Event Ingestion"切成独立 stage,分别打印耗时,最终输出 WorkUnit 数、耗时、峰值内存、usage 去重状态占用磁盘大小与哈希碰撞数。
generate_events的实现细节值得注意:它会按 10% 概率随机把查询"错配"到别的 project(proabability_of_project_mismatch=0.1,模拟审计日志中 project 归属不一致的真实情况),并对视图查询做"下游列访问映射到上游父表"的处理,同时对每个查询生成对应ReadEvent记录fieldsRead。
运行:
python -m tests.performance.bigquery.test_bigquery_usage4.3 Databricks:Unity Catalog 基准 + 真实集群数据灌装
databricks/目录包含两个互补的工具:
test_unity.py:与 Snowflake 测试对称,通过
UnityCatalogApiProxyMock(unity_proxy_mock.py,实现UnityCatalogApiProxy的catalogs()/schemas()/tables()/queries()等接口并内置 schema→table 缓存)替换真实 Databricks SDK,然后用patch("datahub.ingestion.source.unity.source.UnityCatalogApiProxy", lambda *args, **kwargs: proxy_mock)注入。其数据规模为50,000 张表、10,000 个视图、容器层级[1, 100, 5000](1 个 metastore、100 个 catalog、5000 个 schema),每表 100±50 列,外加 20 万条查询与 10,000 个 service principal,配置include_usage_statistics=True。python -m tests.performance.databricks.test_unitygenerator.py:
DatabricksDataGenerator是面向真实 Databricks 集群的灌数器——它用WorkspaceClient+make_sqlalchemy_uri建立连接,把SeedMetadata物化为真实的 catalog/schema/table/view 及行数据(每表行数服从LomaxDistribution(scale=100, shape=1.5),封顶 100 万行),并通过INSERT ... SELECT ... FROM upstream的方式构造真实血缘、用 200 线程池并发执行建表/灌数/建血缘。这为"需要真实执行引擎验证 SQL 语义"的场景(如视图定义正确性)提供了与纯 mock 互补的路径。该文件引入的是performance.xxx包路径(from performance.data_generation import ...),与tests.performance.xxx略有差异,属于进行中的演进痕迹,读者以自身 checkout 版本为准。
4.4 各基准测试规模速览
| 测试 | 数据规模 | 被测对象 | 独特度量 |
|---|---|---|---|
| Snowflake | 30k 表 / 10k 视图 / 30 列 | SnowflakeV2Source(血缘+usage+operational+SQL 格式化) | source.get_report() |
| BigQuery | 20k 表 / 2k 视图 / 104 万查询 | BigQueryUsageExtractor | stage 拆分、usage 状态磁盘占用、哈希碰撞 |
| Databricks (mock) | 50k 表 / 10k 视图 / 20 万查询 | UnityCatalogSource | 多层容器 1/100/5000、service principal |
| Databricks (真实) | 由 SeedMetadata 决定(每表至多 100 万行) | 真实集群 DDL/DML | 线程池 200 并发 |
五、SQL 解析链路的两类专项基准
SQL 解析是 ingestion 中单条记录处理成本最高的环节之一,框架单独提供了两个聚焦测试。
5.1 SQL 格式化吞吐test_sql_formatter.py
sql/test_sql_formatter.py 对datahub.sql_parsing.sqlglot_utils.try_format_query进行 500 次迭代的纯耗时测试,使用来自tests.integration.snowflake.common.large_sql_query的大 SQL 语句,每 50 次打印一次累计耗时。由于被测函数本身带有缓存装饰器,这里刻意调用try_format_query.__wrapped__以绕过缓存、测量真实格式化成本——这是编写微基准时需要记住的技巧。
python -m tests.performance.sql.test_sql_formatter5.2 SQL Aggregator 吞吐回归门禁test_sql_aggregator.py
sql_parsing/test_sql_aggregator.py 与前几个不同,它是一等公民的pytest 性能门禁测试(@pytest.mark.perf),会在吞吐低于阈值时直接断言失败,用于在 CI 中捕获性能回归:
运行方式(文件 docstring 给出了官方命令):
pytest tests/performance/sql_parsing/test_sql_aggregator.py::test_benchmark -s --log-cli-level=INFO规模矩阵:
QUERY_COUNT_OPTIONS = [1000] if is_ci() else [100, 1000, 10000]——CI 下只跑 1000 条(100 条在共享 CI 上测量噪声太大,源码注释说明了这一取舍),本地可跑 100/1000/10000 三档;吞吐阈值:
MIN_THROUGHPUT_THRESHOLD = 50.0 if is_ci() else 90.0,即本地要求 ≥90 queries/sec,CI 放宽到 ≥50;查询生成:
generate_queries_at_scale用三层模板(简单/中等/复杂,含 JOIN、CTE、MERGE INTO、窗口聚合等)随机组合 20 个用户、递增时间戳,seed=42保证可复现(可用SQL_AGGREGATOR_TEST_SEED覆盖);被测对象:
SqlParsingAggregator(platform="redshift", generate_lineage=True, generate_usage_statistics=False, generate_operations=False),逐条add(query)后计时,close()收尾;测试开头会设置DATAHUB_SQL_AGG_SKIP_JOINS=true跳过 join 解析以聚焦核心吞吐;输出:按
query_count / elapsed_time / avg_time_per_query / throughput打印对齐的结果表,并用@pytest.mark.flaky(reruns=5)缓解 CI 抖动。
六、GraphQL 查询投影管线基准test_graphql_projection_perf.py
test_graphql_projection_perf.py 面向 DataHub 客户端侧的一个专门环节:datahub.utilities.graphql_query_adapter中的QueryProjector(把用户 GraphQL 查询裁剪为 GMS 支持的字段子集,涉及parse → _inline_fragments → UnsupportedFieldRemover 访问 → print_ast四阶段)。该基准有两档运行模式:
纯 mock 模式(无需服务器):
pytest tests/performance/test_graphql_projection_perf.py -s -k "not live"连真实 DataHub 实例(对比 introspection 网络往返与冷/热路径):
DATAHUB_PERF_GMS_URL=http://localhost:8080 \ pytest tests/performance/test_graphql_projection_perf.py -s -k live
工作负载(Workload)取自仓库真实产物:CLI 的search.gql/semantic_search.gql(约 735 行、14 个命名 fragment spread)、agent-context 的entity_details.gql(1735 行、93 个 spread)与document_search.gql,外加一个 15 行的最小查询。每个 workload 用内嵌 SDL 构建GraphQLSchema,逐阶段输出median/p95 毫秒数(_time_n重复 100 次取中位数与 95 分位),并额外测量两级缓存的命中成本(Tier-2 dict 命中和 Tier-1 命中+Tier-2 未命中),用于验证"_inline_fragments()的开销可忽略"这一优化结论。live 模式还会强制 TTL 过期触发 schema 重新 introspection,测出 re-fetch 成本。
七、内存泄漏专项:sqlglot[c] 与 SQL 解析缓存
最后一个专项测试 test_sqlglot_memory_leak.py 反映了 DataHub 维护过程中真实踩过的坑:sqlglot的 C 加速实现(sqlglot[c])在反复访问Table.name时存在引用计数泄漏,导致解析大量 BigQuery 视图时内存累积数 GB。该文件用三个@pytest.mark.perf用例把它固化为可回归验证的测试:
test_sqlglot_table_name_memory_leak:解析一条带 JOIN 的 BigQuery 查询,用sys.getrefcount追踪表标识符对象,反复 100 次访问table.name后gc.collect()再比对引用计数,断言增量必须 < 10,否则判定泄漏存在(sqlglot[c]泄漏时该值会接近ITERATIONS × len(tables));test_view_lineage_extraction_memory_usage:用tracemalloc对 160 次sqlglot_lineage(datahub/sql_parsing/sqlglot_lineage.py)调用做快照对比,输出内存增量并按真实环境 16,443 个视图做外推,超过 100MB 时打印告警(该数字对应真实 BigQuery 生产视图规模,来自测试注释);test_parse_cache_memory_footprint:直接检查_sqlglot_lineage_cached这个 LRU 缓存的cache_info()(命中/未命中/当前大小),估算单条SqlParsingResult占用并外推满 1000 条缓存时的总内存。
pytest tests/performance/test_sqlglot_memory_leak.py -s这三个用例展示了性能测试框架的另一层价值:不仅测"快不快",还测"会不会越跑越慢",把偶发的外部依赖问题转变成可重复、可断言的项目资产。
八、如何为新的 ingestion source 接入性能测试
综合前文,把一个新 source 纳入该框架的标准姿势可以归纳为四步:
- 造种子数据:调用
generate_data(num_containers=..., num_tables=..., num_views=..., time_range=...)获得SeedMetadata,必要时用generate_queries补充 usage 查询; - mock 或灌真:对 SDK 型 source(Snowflake/Databricks)用
mock.patch替换连接层并注入按规模放大的default_query_results或自定义 proxy mock;对需要真实引擎验证的场景使用DatabricksDataGenerator这类灌数器; - 跑全管线:构建
PipelineContext(run_id="test")与对应Config,用PerfTimer+workunit_sink包裹source.get_workunits(),记录 WorkUnit 数与峰值内存,打印source.get_report(); - 固化门禁:若需进入 CI 回归体系,参考
test_sql_aggregator.py的模式,用@pytest.mark.perf标记、以吞吐阈值断言(注意区分 CI/本地阈值)、固定随机种子保证可复现,并对测量敏感型用例加@pytest.mark.flaky(reruns=5)容错。
在自行编写时,还应继承前文提到的几条最佳实践:内存采样要节流(如每 10k 个 WorkUnit 一次)、计时用time.perf_counter()、大 SQL 微基准要绕过函数缓存(.__wrapped__)、对长耗用例用 stage 拆分定位瓶颈。
九、小结
DataHub 的metadata-ingestion/tests/performance是一个小而完整的性能测试框架:data_model.py与data_generation.py提供贴近真实幂律分布的合成仓库数据;helpers.py的workunit_sink与datahub.utilities.perf_timer.PerfTimer构成统一的度量口径;Snowflake / BigQuery / Databricks 三个基准覆盖了从全 mock 到真实集群的验证光谱;SQL 格式化、SQL Aggregator 吞吐门禁、GraphQL 投影与内存泄漏测试则把"性能"从单一时延扩展到了吞吐、缓存命中与内存稳定性多个维度。对于任何想要为自己的 ingestion source 建立性能基线的开发者,这个目录都是一个可直接借鉴的蓝本——原 README 中的一行python -m tests.performance.<test_name>背后,是一整套可运行、可度量、可断言的工程实践。
- 数据目录
- 数据治理
- 数据血缘
- 后端
- 前端
- 数据工程
- 数据集成
【免费下载链接】datahub
The Context Platform for your Data and AI Stack
相关推荐
AIBrix Benchmark 基准测试框架指南:从数据集生成到性能分析的端到端配置与实战
AIBrix Benchmark 基准测试框架指南:从数据集生成到性能分析的端到端配置与实战 导读 AIBrix Benchmark 是 AIBrix 项目中用
云原生大模型模型推理服务API网关LLM 网关弹性伸缩可观测性后端POCO内存泄漏检测工具集成:CMake与测试框架
POCO内存泄漏检测工具集成:CMake与测试框架 你还在为内存泄漏头疼?一文解决POCO开发痛点 内存泄漏(Memory Leak)是C++开发中常见的隐患,
后端网络/通信数据库密码学Web框架终极LevelDB测试框架实践指南:从单元测试到性能基准测试的完整教程
终极LevelDB测试框架实践指南:从单元测试到性能基准测试的完整教程 LevelDB是Google开发的一款快速键值存储库,提供从字符串键到字符串值的有序映射
数据库KV存储嵌入式数据库
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考