DataHub Snowplow 连接器本地集成测试环境搭建与验证指南(Option B:DuckDB + Mock BDP Server)
2026/9/20 4:30:56 网站建设 项目流程

DataHub Snowplow 连接器本地集成测试环境搭建与验证指南(Option B:DuckDB + Mock BDP Server)

【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

导读

本文基于 metadata-ingestion/tests/integration/snowplow/docs/SETUP_VERIFICATION.md 展开,完整记录 DataHub 开源仓库中Snowplow 连接器(Option B:本地开发环境)的搭建过程与验证结果:通过增强型 JSON Fixture、DuckDB 本地数仓、Mock BDP API Server 三件套,无需真实 Snowplow BDP 账号即可完成连接器的集成测试、数据血缘(warehouse lineage)与所有权(ownership)抽取验证。读完本文,你将掌握这套本地测试环境的全部组件结构、脚本参数、验证命令与故障排查方法,并能直接在自己的 DataHub 仓库中复现整套流程。


一、Option B 方案概览:为什么需要本地测试环境

Snowplow 连接器的生产目标环境是 Snowplow BDP(Behavioral Data Platform,托管版),其 Console API 需要付费/试用账号,且仓库型数仓依赖外部资源,无法在 CI 中快速、离线地反复验证。为此,仓库在 metadata-ingestion/tests/integration/snowplow 下提供了两套方案:

  • Option A(BDP Cloud):直连真实 BDP 环境,用于上线前的最终验证;
  • Option B(Local):完全本地化——用 JSON Fixture 模拟 API 响应、用 DuckDB 模拟 Snowplow 数仓、用 Flask 起一个 Mock BDP Server 模拟真实 HTTP 调用,约 5 分钟即可搭建,适合开发迭代与 CI/CD。

两者的对比(来源:LOCAL_TEST_SETUP.md):

特性Option A(BDP Cloud)Option B(Local)
搭建耗时1-2 小时5 分钟
外部依赖BDP 账号、Warehouse
成本付费/试用账号免费
真实性生产级 APIMock 响应
离线测试
适用场景最终验证开发与 CI

本文聚焦 Option B,其四个核心组件(增强 Fixtures、DuckDB 脚本、Mock BDP Server、配套文档)均已全部验证通过。


二、组件一:增强型 JSON Fixtures(含所有权/部署历史)

2.1 文件位置与内容

增强后的 Fixture 位于 fixtures/data_structures_with_ownership.json,包含 3 个贴近真实业务的 Snowplow Schema(com.acme命名空间),完整覆盖部署历史、多版本演化与自定义元数据:

  1. checkout_started(event)1-0-0ryan@company.com创建,1-1-0jane@company.com修改,用于测试 Schema 演化与字段作者归属(field authorship);
  2. product_viewed(event):单版本(1-0-0),创建者alice@company.com
  3. user_context(entity):实体(context)类型 Schema,创建者bob@company.com

所有权相关数据要点:

  • deployments数组携带initiator(发起人)字段;
  • ✅ 每个 Schema 有多个版本,用于验证版本历史;
  • ts时间戳支持时间维度追踪;
  • meta.customData中可放置自定义元数据(如{"team": "checkout"})。

2.2 Fixture 与原始文件的对比

仓库中同时保留了原始 Fixture data_structures_response.json 和增强版本,二者差异如下(来源:SETUP_VERIFICATION.md 中的对比表):

特性data_structures_response.jsondata_structures_with_ownership.json
Schema 数量23
含 deployments是(initiator)
所有权数据
Schema 演化单版本多版本
命名空间com.examplecom.acme
用途简单测试所有权/演化完整测试

2.3 源码印证:所有权如何被抽取

增强 Fixture 中的deployments[].initiator是所有权抽取的数据来源。在连接器源码 builders/ownership_builder.py 中,extract_ownership_from_deployments()按时间戳排序:

  • 最早(最旧)的 deployment →createdBy,映射为 DataHub 的DATAOWNER所有权类型;
  • 最新(最近)的 deployment →modifiedBy,映射为PRODUCER
  • 字段级作者(field authorship)则由版本历史推导而来。

同时在 test_snowplow.py 中,test_snowplow_ingest会加载该 Fixture,mockSnowplowBDPClient.get_data_structures,并 mock/users返回的 4 个用户(ryan@company.comjane@company.comalice@company.combob@company.com)用于将initiator解析为真实用户 URN,最终与 golden 文件比对。


三、组件二:DuckDB 本地数仓脚本

3.1 脚本能力

脚本位于 setup/setup_duckdb.py,用于在本地创建模拟 Snowplowatomic.events表的 DuckDB 数据库:

  • 创建本地 DuckDB 数据库;
  • 生成贴近真实的snowplow.events表;
  • 灌入示例事件数据;
  • 支持自定义事件数量与数据库路径。

3.2 命令行参数

参数默认值说明
--db-pathsnowplow_test.duckdbDuckDB 数据库文件路径
--event-count100生成的示例事件数量
--recreate关闭删除并重建数据库

(参数定义见 setup_duckdb.py 的 argparse 部分。)

3.3 数据库表结构

create_database()(setup_duckdb.py)创建的snowplow.events表结构如下:

snowplow.events ( -- 基础列(Base columns) app_id, platform, collector_tstamp, event, event_id, user_id, user_ipaddress, -- 页面上下文(Page context) page_url, page_title, page_referrer, -- 设备上下文(Device context) br_name, br_family, os_name, os_family, -- Geo 富化(Geo enrichment) geo_country, geo_region, geo_city, geo_zipcode, geo_latitude, geo_longitude, -- 自定义事件上下文(JSON 列) contexts_com_acme_checkout_started_1, contexts_com_acme_product_viewed_1, contexts_com_acme_user_context_1, -- Unstruct 事件(JSON 列) unstruct_event_com_acme_checkout_started_1, unstruct_event_com_acme_product_viewed_1, -- 时间戳 derived_tstamp, load_tstamp )

3.4 样例数据生成逻辑

generate_sample_events()(setup_duckdb.py)在 30 天时间窗内以 5 分钟为间隔生成事件:

  • 随机选择事件类型:checkout_started(映射为unstruct事件)或product_viewed(映射为struct事件);
  • 每行生成event_id(UUID)、user_iduser_1~user_100)、随机浏览器/操作系统/地理位置;
  • checkout_started的 unstruct/context JSON 含amountcurrency(USD/EUR/GBP)、discount_codeSAVE10/WELCOME20)、items数组;
  • product_viewed的 JSON 含product_idcategory(electronics/clothing/books)、price
  • 每行附加contexts_com_acme_user_context_1user_type取值 free/premium/enterprise、registration_date)。

当前预置数据库 setup/snowplow_test.duckdb 共 100 条事件,分布为struct 54 条 / unstruct 46 条,可通过分组查询验证:

SELECT event, COUNT(*) FROM snowplow.events GROUP BY event

输出:

eventcount
struct54
unstruct46

四、组件三:Mock BDP API Server

4.1 脚本能力与启动方式

脚本位于 setup/mock_bdp_server.py,基于 Flask 实现,模拟 Snowplow BDP Console API,从 Fixture 目录读取数据返回,支持过滤、分页、完整请求/响应日志与错误处理。启动命令:

python mock_bdp_server.py --port 8081 # 可选参数:--host(默认 localhost)、--debug(开启调试模式)

4.2 端点清单

对照实际源码(mock_bdp_server.py),Mock Server 提供的端点如下:

端点说明源码位置
GET /API 文档与端点索引L326-L346
GET /health健康检查L314-L323
GET /organizations/{orgId}/credentials/v3/tokenJWT Token 签发L68-L93
GET /organizations/{orgId}/data-structures/v1Schema 列表(支持 filter/vendor/limit/offset)L96-L147
GET /organizations/{orgId}/data-structures/v1/{hash}按 hash 获取单个 SchemaL150-L174
GET /organizations/{orgId}/data-products/v2Data products(实际代码为 v2 端点)L177-L204
GET /organizations/{orgId}/users组织用户列表(用于 initiator 邮箱解析)L207-L250
GET /organizations/{orgId}/event-specs/v1Event specificationsL253-L280
GET /organizations/{orgId}/tracking-scenarios/v1Tracking scenarios(旧路径)L283-L311

说明:原验证文档中列出的是/data-products/v1,而当前仓库源码实现为/data-products/v2(v2 使用data/includes/errors包裹格式,与 event-specs 一致,区别于>source: type: snowplow config: # BDP 连接(Mock——需先启动 mock server) bdp_connection: organization_id: "test-org-uuid" api_key_id: "test-key-id" api_key: "test-secret" console_api_url: "http://localhost:8081/api/msc/v1" # Mock server # Schema 过滤 schema_pattern: allow: - "com.acme.*" # 只抽取 event/entity 类型 Schema schema_types_to_extract: - "event" - "entity" # 可选功能——本测试中关闭 extract_event_specifications: false extract_tracking_plans: false # 仓库血缘——通过 DuckDB 启用 extract_warehouse_lineage: true # DuckDB 仓库连接 warehouse_connection: warehouse_type: "duckdb" database: "snowplow_test.duckdb" schema_name: "snowplow" # 平台标识 platform_instance: "snowplow-test" env: "TEST" sink: type: file config: filename: "./snowplow_duckdb_output.json"

其中console_api_url指向 Mock Server,使连接器通过真实 HTTP 链路访问本地模拟 API。相关连接参数(organization_idapi_key_idapi_keyconsole_api_url默认值https://console.snowplowanalytics.com/api/msc/v1timeout_seconds默认 60、max_retries默认 3)在源码 snowplow_config.py 的SnowplowBDPConnectionConfig中定义。


六、验证测试记录与命令

6.1 Test 1:Mock API 集成测试 ✅ PASSED

cd metadata-ingestion source venv/bin/activate pytest tests/integration/snowplow/test_snowplow.py::test_snowplow_ingest -v

验证点:Mock API 响应正常、golden 文件比对通过、无错误。该测试通过unittest.mock.patch替换datahub.ingestion.source.snowplow.snowplow.SnowplowBDPClient,将内存中的DataStructure对象注入客户端,避免任何真实 HTTP 调用(见 test_snowplow.py)。

6.2 Test 2:DuckDB 数据库搭建 ✅ COMPLETED

cd metadata-ingestion/tests/integration/snowplow/setup python setup_duckdb.py --event-count 100

结果:数据库snowplow_test.duckdb创建成功,Schemasnowplow、表snowplow.events建立,插入 100 条事件。验证查询:

SELECT event, COUNT(*) FROM snowplow.events GROUP BY event

6.3 Test 3:Mock BDP Server ✅

python mock_bdp_server.py --port 8081

逐端点验证:

  1. 健康检查curl http://localhost:8081/healthhealthy

  2. Token 签发

    curl -H "X-Api-Key-Id: test-key-id" -H "X-Api-Key: test-secret" \ http://localhost:8081/organizations/test-org-uuid/credentials/v3/token

    mock_jwt_token_12345

  3. Data Structures 列表

    curl -H "Authorization: Bearer mock_jwt_token_12345" \ http://localhost:8081/organizations/test-org-uuid/data-structures/v1

    → 返回 3 个带 deployments 的 Schema。

6.4 Test 4:Fixture 文件校验 ✅

  • data_structures_response.json(原始):2 个 Schema(page_view、user_context),无 deployment 数据,结构简单,JSON 有效;
  • data_structures_with_ownership.json(增强):3 个 Schema,完整部署历史、initiator 所有权、多版本演化,JSON 有效。

七、快速上手命令速查

以下命令均以仓库根目录为基准:

运行全部集成测试

cd metadata-ingestion source venv/bin/activate pytest tests/integration/snowplow/ -v

初始化 DuckDB(仓库测试用)

cd metadata-ingestion/tests/integration/snowplow/setup python setup_duckdb.py --event-count 100

启动 Mock BDP Server

cd metadata-ingestion/tests/integration/snowplow/setup python mock_bdp_server.py --port 8081

查询 DuckDB

cd metadata-ingestion/tests/integration/snowplow/setup duckdb snowplow_test.duckdb "SELECT COUNT(*) FROM snowplow.events"

更多查询示例(来自 LOCAL_TEST_SETUP.md):

# 查看 checkout 事件 duckdb snowplow_test.duckdb "SELECT event_id, user_id, unstruct_event_com_acme_checkout_started_1 FROM snowplow.events WHERE event = 'unstruct' LIMIT 5" # 查看用户上下文 duckdb snowplow_test.duckdb "SELECT user_id, contexts_com_acme_user_context_1 FROM snowplow.events LIMIT 5"

八、三类典型测试场景

场景 1:基础集成测试(最快)

目的:验证连接器在 Mock API 响应下正常工作。

pytest metadata-ingestion/tests/integration/snowplow/test_snowplow.py -v

状态:✅ Working。适合开发期的快速校验,无需任何外部依赖。

场景 2:DuckDB 仓库血缘(warehouse lineage)

目的:从本地 DuckDB 验证仓库血缘抽取。

前置条件

  1. 已创建 DuckDB 数据库(python setup_duckdb.py --event-count 100);
  2. Mock BDP Server 已运行。

步骤

# 终端 1:启动 mock server python metadata-ingestion/tests/integration/snowplow/setup/mock_bdp_server.py --port 8081 # 终端 2:执行 ingestion datahub ingest -c metadata-ingestion/tests/integration/snowplow/recipes/snowplow_with_duckdb.yml

状态:🟡 Ready to test(需 Mock Server 与 DuckDB 联动)。仓库血缘的抽取逻辑可进一步参考连接器源码 processors/warehouse_lineage_processor.py。

场景 3:完整 Mock 环境(端到端)

目的:通过真实 HTTP 调用链路做端到端测试。

组件:Mock BDP Server(HTTP)+ DuckDB 数据库(数仓)+ 测试 Recipe。

# 终端 1:启动 mock server cd metadata-ingestion/tests/integration/snowplow/setup python mock_bdp_server.py --port 8081 # 终端 2:初始化 DuckDB python setup_duckdb.py --event-count 100 # 终端 3:配置环境变量并执行 ingestion export SNOWPLOW_ORG_ID="test-org-uuid" export SNOWPLOW_API_KEY_ID="test-key-id" export SNOWPLOW_API_KEY="test-secret" datahub ingest -c metadata-ingestion/tests/integration/snowplow/recipes/snowplow_with_duckdb.yml

状态:🟡 基础设施就绪,待完整端到端验证。


九、目录结构速览

metadata-ingestion/tests/integration/snowplow/ ├── fixtures/ │ ├── data_structures_response.json [原始 Fixture] │ └── data_structures_with_ownership.json [增强:含 deployments] ├── golden_files/ │ └── snowplow_mces_golden.json [Golden 比对文件] ├── recipes/ │ ├── snowplow_with_duckdb.yml [DuckDB 数仓 Recipe] │ └── test_mock_bdp.yml 等 [其他测试 Recipe] ├── setup/ │ ├── setup_duckdb.py [数据库搭建脚本] │ ├── mock_bdp_server.py [Mock API Server] │ └── snowplow_test.duckdb [预置测试数据库] ├── docs/ │ ├── LOCAL_TEST_SETUP.md [本地使用指南] │ └── SETUP_VERIFICATION.md [验证报告(本文依据)] ├── test_snowplow.py [集成测试] └── test_snowplow_performance.py [性能测试]

完整结构与测试分类可参考 tests/integration/snowplow/README.md。


十、依赖安装

  • datahub[snowplow]:带 Snowplow 支持的主包;
  • duckdb:本地数仓测试;
  • flask:Mock BDP Server(可选)。

安装方式:

cd metadata-ingestion source venv/bin/activate pip install -e ".[snowplow]" duckdb flask

十一、后续工作项(Next Steps)

已完成的即时项

  1. ✅ 基础测试通过;
  2. ✅ DuckDB 数据库已创建;
  3. ✅ Mock Server 已验证;
  4. ✅ Fixtures 校验通过。

待完成的完整测试项

  1. 🟡 用 DuckDB 测试仓库血缘抽取;
  2. 🟡 从 deployments 测试所有权抽取(可参考 docs/OWNERSHIP_TESTING_GUIDE.md);
  3. 🟡 用多版本测试 Schema 演化;
  4. 🟡 测试 data products 抽取(启用后)。

连接器对接项

  1. 更新连接器以支持 DuckDB 仓库类型;
  2. 实现从 deployments 的所有权抽取(当前OwnershipBuilder.extract_ownership_from_deployments已具备基础能力,见 builders/ownership_builder.py);
  3. 从版本历史增加字段作者归属;
  4. 完成全链路端到端测试。

十二、常见问题排查(Troubleshooting)

问题:Module not found 报错

解决:确认 venv 已激活且包已安装:

cd metadata-ingestion source venv/bin/activate pip install -e ".[snowplow]"

问题:DuckDB 未找到

解决:安装 duckdb:

pip install duckdb

问题:Mock Server 端口被占用

解决:更换端口或停止占用进程:

lsof -i :8081 # 查找占用进程 kill <PID> # 停止进程 # 或换端口: python mock_bdp_server.py --port 8082

问题:数据库文件不存在

解决:重新创建数据库:

cd metadata-ingestion/tests/integration/snowplow/setup python setup_duckdb.py

问题:集成测试失败

解决:若输出属预期变更,可更新 golden 文件:

pytest tests/integration/snowplow/test_snowplow.py --update-golden-files # 或查看详细 diff pytest tests/integration/snowplow/test_snowplow.py -vv

问题:Fixture 文件缺失

解决:确认文件存在并从 git 恢复:

ls -la metadata-ingestion/tests/integration/snowplow/fixtures/ git checkout metadata-ingestion/tests/integration/snowplow/fixtures/data_structures_with_ownership.json

十三、自定义测试数据

新增 Schema Fixture

编辑 data_structures_with_ownership.json,按以下模板追加:

{ "hash": "new_schema_hash", "organizationId": "test-org-uuid", "vendor": "com.acme", "name": "new_event", "format": "jsonschema", "description": "New event schema", "meta": { "hidden": false, "schemaType": "event", "customData": { "team": "new-team" } }, "deployments": [ { "version": "1-0-0", "initiator": "developer@company.com", "ts": "2024-03-01T10:00:00Z" } ], "data": { "self": { "vendor": "com.acme", "name": "new_event", "version": "1-0-0" }, "properties": { "field1": { "type": "string" } } } }

调整 DuckDB 事件数据

修改 setup/setup_duckdb.py 中的generate_sample_events():可新增事件类型、调整字段分布、改变事件数量或时间范围。

扩展 Mock API 端点

修改 setup/mock_bdp_server.py:新增 Flask 路由、创建新的 Fixture 文件或实现新的 API 行为。


十四、结论

Option B 本地测试环境已完整搭建并通过验证

  • 集成测试全部通过(golden 文件比对无差异);
  • DuckDB 数据库创建成功,含 100 条示例事件(struct 54 / unstruct 46);
  • Mock BDP Server 正常提供增强 Fixture;
  • 配套文档齐全。

该环境已具备支撑以下工作的能力:

  • ✅ 快速开发迭代(无外部依赖、5 分钟可复现);
  • ✅ 离线测试;
  • ✅ CI/CD 集成;
  • ✅ 基于 deployments 的所有权测试;
  • ✅ 基于 DuckDB 的仓库血缘测试。

推荐后续动作

  1. 使用 DuckDB 验证仓库血缘抽取;
  2. 实现从 deployments 数组的所有权抽取;
  3. 结合 Mock Server 完成端到端全链路测试;
  4. 让连接器集成测试切换到增强 Fixture。

如需进行真实 BDP 环境(Option A)的最终上线验证,可参考 docs/REAL_BDP_TESTING_GUIDE.md。

【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询