
数据目录数据治理数据血缘后端前端数据工程数据集成【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址https://gitcode.com/GitHub_Trending/da/datahub点击查看免费下载本篇技术指南以 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_sinkhelpers.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(digits2)输出耗时。其实现位于 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_sizeusage 去重状态占用磁盘大小与num_usage_query_hash_collisions查询哈希碰撞数。三、公共组件数据模型与合成数据生成器仓库型 source 的性能测试不能依赖真实生产环境既不安全也不可重复因此框架先用程序化生成的方式构造一份接近真实的元数据宇宙再喂给被测 source。3.1 数据模型data_model.pydata_model.py 定义了用于构建模拟仓库的最小领域模型Container层级容器如 database → schema支持parent指针形成多层嵌套Column/ColumnType列与列类型枚举INTEGER、FLOAT、STRING、BOOLEAN、DATETIME基于StrEnum实现Table/View表与视图Table.columns用OrderedDict保持列序upstreams记录血缘上游View额外持有definition视图定义 SQLname_components属性沿容器链展开得到[catalog, schema, table]形式的全限定名这与 BigQuery/Databricks 的标识符模型完全一致FieldAccess一次查询中哪个表上的哪一列被访问用于 usage 场景Query一条审计日志级别的查询记录字段包括textSQL 文本、typeStatementType取值SELECT/INSERT/UPDATE/DELETE/CREATE/ALTER/DROP/CUSTOM/UNKNOWN、actor执行用户、timestamp、fields_accessed与可选的object_modified被修改的对象。3.2 分布与数据生成data_generation.pydata_generation.py 是框架最核心的公共工具它并不生成平均分布的死板数据而是让数据特征贴近真实仓库的幂律形态Distribution抽象 两种实现NormalDistribution(mu, sigma)正态分布适合列数、查询长度等集中在均值附近的指标LomaxDistribution(scale, shape)重尾分布等价于pareto(scale, shape) - scale适合模拟真实环境中的少数热表被大量查询、极少数超大上游现象——源码注释明确给出了该分布在血缘上游数量上的百分位形态75th0、80th1、95th2、99th4、99.99th15。二者均可通过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_rangeusage 时间窗默认 14 天。 它返回SeedMetadatacontainers分层容器列表、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_tables30000, num_views10000, num_cols30, num_ops30, num_usages500, )default_query_results来自 tests/integration/snowflake/common.py是集成测试共享的模拟查询结果集工厂这里把它放大到30,000 张表、10,000 个视图、每表 30 列、30 条操作、500 条 usage的量级SnowflakeV2Config配置了include_technical_schemaFalse、include_table_lineageTrue、include_usage_statsTrue、include_operational_statsTrue、format_sql_queriesTrue即同时压测技术元数据、血缘、usage 与 SQL 格式化四条链路计时与度量沿用前文的PerfTimerworkunit_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_containers2000, num_tables20000, num_views2000, time_rangetimedelta(days7))生成 2 万表、2 千视图、2 千容器事件生成100 个项目随机分配给表随后generate_queries(..., num_selects240_000, num_operations800_000, num_unique_queries50_000, num_users2000, query_lengthNormalDistribution(2000, 500))产生104 万条查询按时间排序后经 bigquery_events.py 的generate_events转成AuditEventQueryEventReadEvent流被测对象BigQueryUsageExtractor来自 datahub/ingestion/source/bigquery_v2/usage.py配置usage.include_top_n_queriesTrue、top_n_queries10、apply_view_usage_to_tablesTrue并设置file_backed_cache_size1000文件后备缓存上限分阶段度量BigQueryV2Report.new_stage(...)把Seed Data Generation / Event Generation / Event Ingestion切成独立 stage分别打印耗时最终输出 WorkUnit 数、耗时、峰值内存、usage 去重状态占用磁盘大小与哈希碰撞数。generate_events的实现细节值得注意它会按 10% 概率随机把查询错配到别的 projectproabability_of_project_mismatch0.1模拟审计日志中 project 归属不一致的真实情况并对视图查询做下游列访问映射到上游父表的处理同时对每个查询生成对应ReadEvent记录fieldsRead。运行python -m tests.performance.bigquery.test_bigquery_usage4.3 DatabricksUnity Catalog 基准 真实集群数据灌装databricks/目录包含两个互补的工具test_unity.py与 Snowflake 测试对称通过UnityCatalogApiProxyMockunity_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_statisticsTrue。python -m tests.performance.databricks.test_unitygenerator.pyDatabricksDataGenerator是面向真实 Databricks 集群的灌数器——它用WorkspaceClientmake_sqlalchemy_uri建立连接把SeedMetadata物化为真实的 catalog/schema/table/view 及行数据每表行数服从LomaxDistribution(scale100, shape1.5)封顶 100 万行并通过INSERT ... SELECT ... FROM upstream的方式构造真实血缘、用 200 线程池并发执行建表/灌数/建血缘。这为需要真实执行引擎验证 SQL 语义的场景如视图定义正确性提供了与纯 mock 互补的路径。该文件引入的是performance.xxx包路径from performance.data_generation import ...与tests.performance.xxx略有差异属于进行中的演进痕迹读者以自身 checkout 版本为准。4.4 各基准测试规模速览测试数据规模被测对象独特度量Snowflake30k 表 / 10k 视图 / 30 列SnowflakeV2Source血缘usageoperationalSQL 格式化source.get_report()BigQuery20k 表 / 2k 视图 / 104 万查询BigQueryUsageExtractorstage 拆分、usage 状态磁盘占用、哈希碰撞Databricks (mock)50k 表 / 10k 视图 / 20 万查询UnityCatalogSource多层容器 1/100/5000、service principalDatabricks (真实)由 SeedMetadata 决定每表至多 100 万行真实集群 DDL/DML线程池 200 并发五、SQL 解析链路的两类专项基准SQL 解析是 ingestion 中单条记录处理成本最高的环节之一框架单独提供了两个聚焦测试。5.1 SQL 格式化吞吐test_sql_formatter.pysql/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.pysql_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-levelINFO规模矩阵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/secCI 放宽到 ≥50查询生成generate_queries_at_scale用三层模板简单/中等/复杂含 JOIN、CTE、MERGE INTO、窗口聚合等随机组合 20 个用户、递增时间戳seed42保证可复现可用SQL_AGGREGATOR_TEST_SEED覆盖被测对象SqlParsingAggregator(platformredshift, generate_lineageTrue, generate_usage_statisticsFalse, generate_operationsFalse)逐条add(query)后计时close()收尾测试开头会设置DATAHUB_SQL_AGG_SKIP_JOINStrue跳过 join 解析以聚焦核心吞吐输出按query_count / elapsed_time / avg_time_per_query / throughput打印对齐的结果表并用pytest.mark.flaky(reruns5)缓解 CI 抖动。六、GraphQL 查询投影管线基准test_graphql_projection_perf.pytest_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_URLhttp://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.gql1735 行、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_lineagedatahub/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 型 sourceSnowflake/Databricks用mock.patch替换连接层并注入按规模放大的default_query_results或自定义 proxy mock对需要真实引擎验证的场景使用DatabricksDataGenerator这类灌数器跑全管线构建PipelineContext(run_idtest)与对应Config用PerfTimerworkunit_sink包裹source.get_workunits()记录 WorkUnit 数与峰值内存打印source.get_report()固化门禁若需进入 CI 回归体系参考test_sql_aggregator.py的模式用pytest.mark.perf标记、以吞吐阈值断言注意区分 CI/本地阈值、固定随机种子保证可复现并对测量敏感型用例加pytest.mark.flaky(reruns5)容错。在自行编写时还应继承前文提到的几条最佳实践内存采样要节流如每 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背后是一整套可运行、可度量、可断言的工程实践。赞分享数据目录数据治理数据血缘后端前端数据工程数据集成【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址https://gitcode.com/GitHub_Trending/da/datahub点击查看免费下载相关推荐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),仅供参考