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 | 无 |
| 成本 | 付费/试用账号 | 免费 |
| 真实性 | 生产级 API | Mock 响应 |
| 离线测试 | 否 | 是 |
| 适用场景 | 最终验证 | 开发与 CI |
本文聚焦 Option B,其四个核心组件(增强 Fixtures、DuckDB 脚本、Mock BDP Server、配套文档)均已全部验证通过。
二、组件一:增强型 JSON Fixtures(含所有权/部署历史)
2.1 文件位置与内容
增强后的 Fixture 位于 fixtures/data_structures_with_ownership.json,包含 3 个贴近真实业务的 Snowplow Schema(com.acme命名空间),完整覆盖部署历史、多版本演化与自定义元数据:
- checkout_started(event):
1-0-0由ryan@company.com创建,1-1-0由jane@company.com修改,用于测试 Schema 演化与字段作者归属(field authorship); - product_viewed(event):单版本(
1-0-0),创建者alice@company.com; - 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.json | data_structures_with_ownership.json |
|---|---|---|
| Schema 数量 | 2 | 3 |
| 含 deployments | 否 | 是(initiator) |
| 所有权数据 | 否 | 是 |
| Schema 演化 | 单版本 | 多版本 |
| 命名空间 | com.example | com.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.com、jane@company.com、alice@company.com、bob@company.com)用于将initiator解析为真实用户 URN,最终与 golden 文件比对。
三、组件二:DuckDB 本地数仓脚本
3.1 脚本能力
脚本位于 setup/setup_duckdb.py,用于在本地创建模拟 Snowplowatomic.events表的 DuckDB 数据库:
- 创建本地 DuckDB 数据库;
- 生成贴近真实的
snowplow.events表; - 灌入示例事件数据;
- 支持自定义事件数量与数据库路径。
3.2 命令行参数
| 参数 | 默认值 | 说明 |
|---|---|---|
--db-path | snowplow_test.duckdb | DuckDB 数据库文件路径 |
--event-count | 100 | 生成的示例事件数量 |
--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_id(user_1~user_100)、随机浏览器/操作系统/地理位置; checkout_started的 unstruct/context JSON 含amount、currency(USD/EUR/GBP)、discount_code(SAVE10/WELCOME20)、items数组;product_viewed的 JSON 含product_id、category(electronics/clothing/books)、price;- 每行附加
contexts_com_acme_user_context_1(user_type取值 free/premium/enterprise、registration_date)。
当前预置数据库 setup/snowplow_test.duckdb 共 100 条事件,分布为struct 54 条 / unstruct 46 条,可通过分组查询验证:
SELECT event, COUNT(*) FROM snowplow.events GROUP BY event输出:
| event | count |
|---|---|
| struct | 54 |
| unstruct | 46 |
四、组件三: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/token | JWT Token 签发 | L68-L93 |
GET /organizations/{orgId}/data-structures/v1 | Schema 列表(支持 filter/vendor/limit/offset) | L96-L147 |
GET /organizations/{orgId}/data-structures/v1/{hash} | 按 hash 获取单个 Schema | L150-L174 |
GET /organizations/{orgId}/data-products/v2 | Data products(实际代码为 v2 端点) | L177-L204 |
GET /organizations/{orgId}/users | 组织用户列表(用于 initiator 邮箱解析) | L207-L250 |
GET /organizations/{orgId}/event-specs/v1 | Event specifications | L253-L280 |
GET /organizations/{orgId}/tracking-scenarios/v1 | Tracking 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_id、api_key_id、api_key、console_api_url默认值https://console.snowplowanalytics.com/api/msc/v1、timeout_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 event6.3 Test 3:Mock BDP Server ✅
python mock_bdp_server.py --port 8081逐端点验证:
健康检查:
curl http://localhost:8081/health→healthy;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;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 验证仓库血缘抽取。
前置条件:
- 已创建 DuckDB 数据库(
python setup_duckdb.py --event-count 100);- 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)
已完成的即时项
- ✅ 基础测试通过;
- ✅ DuckDB 数据库已创建;
- ✅ Mock Server 已验证;
- ✅ Fixtures 校验通过。
待完成的完整测试项
- 🟡 用 DuckDB 测试仓库血缘抽取;
- 🟡 从 deployments 测试所有权抽取(可参考 docs/OWNERSHIP_TESTING_GUIDE.md);
- 🟡 用多版本测试 Schema 演化;
- 🟡 测试 data products 抽取(启用后)。
连接器对接项
- 更新连接器以支持 DuckDB 仓库类型;
- 实现从 deployments 的所有权抽取(当前
OwnershipBuilder.extract_ownership_from_deployments已具备基础能力,见 builders/ownership_builder.py);- 从版本历史增加字段作者归属;
- 完成全链路端到端测试。
十二、常见问题排查(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 的仓库血缘测试。
推荐后续动作:
- 使用 DuckDB 验证仓库血缘抽取;
- 实现从 deployments 数组的所有权抽取;
- 结合 Mock Server 完成端到端全链路测试;
- 让连接器集成测试切换到增强 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),仅供参考