DataHub Snowplow 连接器本地测试环境搭建指南:基于 Mock BDP API 与 DuckDB 的零依赖集成测试方案
DataHub Snowplow 连接器本地测试环境搭建指南基于 Mock BDP API 与 DuckDB 的零依赖集成测试方案【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub导读本文基于 DataHub 仓库中 Snowplow 集成测试目录metadata-ingestion/tests/integration/snowplow的本地测试方案Option B系统讲解如何在没有 Snowplow BDP 商业账号的前提下通过 JSON Fixture 模拟 API 响应、DuckDB 模拟数据仓库、Flask Mock 服务器模拟 BDP Console API完成从单元级 Mock 测试到端到端 HTTP 摄入ingestion的完整本地验证闭环。读完本文你将掌握 Snowplow 连接器集成测试的全部本地化手段并能据此搭建属于自己的离线、确定性、可进 CI 的测试环境。一、背景为什么需要本地测试环境Option BSnowplow 连接器在真实环境中依赖 Snowplow BDPBehavioral Data PlatformConsole API 与数据仓库来完成元数据抽取schema、数据产品、管线、事件规格与血缘提取。但这带来三个现实问题需要付费/试用 BDP 账号无法在 CI 中稳定复现依赖真实仓库测试数据不可控、不可重复无法离线开发且网络抖动会引入不确定性。为此该集成测试目录提供了两套并行方案详见对比表而 LOCAL_TEST_SETUP.md 正是 Option B本地方案的操作指南FeatureOption A (BDP Cloud)Option B (Local)Setup Time1-2 hours5 minutesExternal DependenciesBDP account, WarehouseNoneCostRequires paid/trial accountFreeRealismProduction-like APIMocked responsesOffline Testing❌ No✅ YesBest ForFinal validationDevelopment CIOption B 的核心思路是用三件本地组件替代真实依赖Mock API responsesfixtures/目录下的 JSON 文件模拟 BDP Console API 的返回DuckDB 数据库模拟数据仓库中的snowplow.events原子事件表可选 Mock BDP 服务器基于 Flask 的本地 HTTP 服务完整模拟 BDP 的鉴权与数据接口用于带真实 HTTP 调用的端到端测试。二、目录结构与文件角色实际仓库中的目录结构如下注意相比文档中的示意图脚本实际位于setup/子目录recipe 位于recipes/子目录metadata-ingestion/tests/integration/snowplow/ ├── test_snowplow.py # 主集成测试golden file 比对 ├── test_snowplow_performance.py # 性能测试并行抓取、缓存优化 ├── fixtures/ # Mock API 响应 Fixture │ ├── data_structures_response.json # 基础 schema fixture │ ├── data_structures_with_ownership.json # 含 deployments/ownership 的增强 fixture │ ├── data_structures_response_real.json │ ├── enrichments_response.json │ ├── event_specifications_response.json │ ├── organization_response.json │ ├── pipelines_response.json │ ├── tracking_plans_response.json │ └── tracking_scenarios_response.json ├── golden_files/ # 期望输出golden file │ ├── snowplow_mces_golden.json │ ├── snowplow_event_specs_golden.json │ ├── snowplow_tracking_plans_golden.json │ ├── snowplow_pipelines_golden.json │ ├── snowplow_enrichments_golden.json │ └── snowplow_iglu_autodiscovery_golden.json ├── recipes/ # 测试用摄入 recipe │ ├── test_mock_bdp.yml # 纯 Mock BDP 服务器摄入 │ ├── snowplow_with_duckdb.yml # DuckDB 仓库血缘摄入 │ ├── test_ownership_recipe.yml │ ├── test_event_specs.yml │ ├── test_iglu_autodiscovery.yml │ ├── test_real_bdp.yml / test_real_bdp_to_datahub.yml / test_real_bdp_to_file.yml │ └── test_datahub_ui.yml ├── setup/ │ ├── mock_bdp_server.py # Mock BDP API 服务器Flask │ ├── setup_duckdb.py # DuckDB 测试库生成脚本 │ ├── setup_iglu.py # Iglu Server 初始化脚本 │ ├── docker-compose.iglu.yml # Iglu Server 的 Docker Compose │ ├── iglu_config.hocon │ └── snowplow_test.duckdb # 预置好的测试数据库 └── docs/ # 配套文档 ├── LOCAL_TEST_SETUP.md # 本文对应的本地测试指南 ├── REAL_BDP_TESTING_GUIDE.md # Option A真实 BDP 测试指南 ├── OWNERSHIP_TESTING_GUIDE.md # 所有权提取测试指南 ├── SETUP_VERIFICATION.md # 环境验证清单 ├── BDP_API_VALIDATION.md └── SWAGGER_VALIDATION_REPORT.md三、快速开始三种由浅入深的验证方式方式 1纯 Mock 单元测试最快零依赖这是最简单的入门路径——直接运行集成测试测试内部通过unittest.mock.patch拦截SnowplowBDPClient从fixtures/读取模拟响应无需任何外部服务# 在 metadata-ingestion 目录下执行 pytest tests/integration/snowplow/test_snowplow.py -v从 test_snowplow.py 源码可见其工作方式用time_machine.travel(2024-01-01 00:00:00)冻结系统时间保证输出可重复patch(datahub.ingestion.source.snowplow.snowplow.SnowplowBDPClient)替换真实客户端将 fixture 解析为DataStructure/User模型对象后作为 mock 返回值以Pipeline.create({...})驱动摄入输出到临时文件最后调用mce_helpers.check_golden_file(...)与golden_files/snowplow_mces_golden.json比对忽略时间戳、lastObserved、runId等不稳定字段。适用场景开发迭代期的快速验证。方式 2DuckDB 仓库测试验证血缘提取当需要验证「Snowplow 事件表 → DataHub 血缘」链路时先准备一个本地 DuckDB 数据仓库cd metadata-ingestion/tests/integration/snowplow/setup # 创建带 100 条示例事件的数据库 python setup_duckdb.py --event-count 100 # 验证数据库内容 duckdb snowplow_test.duckdb SELECT COUNT(*) FROM snowplow.events随后使用 snowplow_with_duckdb.yml 执行摄入datahub ingest -c ../recipes/snowplow_with_duckdb.yml该 recipe 开启了extract_warehouse_lineage: true并配置了warehouse_connectionwarehouse_type: duckdb、database: snowplow_test.duckdb、schema_name: snowplow同时支持schema_pattern过滤与schema_types_to_extractevent/entity选择输出到文件 sink 便于检查。适用场景验证 warehouse lineage 提取逻辑。方式 3Mock BDP 服务器端到端测试最完整启动本地 Mock 服务器后连接器会走真实的 HTTP 调用链鉴权 → 拉取 schema → 解析 → 落盘# 终端 1启动 Mock 服务器位于 setup 目录 cd metadata-ingestion/tests/integration/snowplow/setup python mock_bdp_server.py --port 8081 # 终端 2设置连接凭证并执行摄入 export SNOWPLOW_ORG_IDtest-org-uuid export SNOWPLOW_API_KEY_IDtest-key-id export SNOWPLOW_API_KEYtest-secret datahub ingest -c ../recipes/test_mock_bdp.ymltest_mock_bdp.yml 中的连接配置为source: type: snowplow config: bdp_connection: organization_id: test-org-uuid api_key_id: test-key-id api_key: test-secret console_api_url: http://localhost:8081 # Mock BDP server extract_event_specifications: false extract_tracking_plans: false sink: type: file config: filename: /tmp/snowplow_ownership_mock_test.json适用场景带真实 HTTP 的端到端验证、错误场景与边界条件测试。四、深入 DuckDB 仓库搭建4.1 脚本参数setup_duckdb.py 暴露三个命令行参数参数默认值说明--db-pathsnowplow_test.duckdb数据库文件路径--event-count100生成的示例事件条数--recreate关闭删除并重建数据库用于重置脏数据示例python setup_duckdb.py --db-path /tmp/my_snowplow.duckdb --event-count 500 --recreate4.2 数据库 Schema脚本会创建snowplow.events表模拟 Snowplow 仓库中的atomic.events结构包含六类字段源码中完整定义这里按类型归纳snowplow.events ( -- 基础 Snowplow 列 app_id, platform, etl_tstamp, collector_tstamp, dvce_created_tstamp, event, event_id, txn_id, name_tracker, -- 用户标识 user_id, user_ipaddress, -- 页面上下文 page_url, page_title, page_referrer, -- 设备上下文 br_name, br_family, os_name, os_family, -- Geo 上下文IP Lookup enrichment 产物 geo_country, geo_region, geo_city, geo_zipcode, geo_latitude, geo_longitude, -- 自定义事件上下文JSON 列 contexts_com_acme_checkout_started_1 JSON, contexts_com_acme_product_viewed_1 JSON, contexts_com_acme_user_context_1 JSON, -- 非结构化事件self-describing eventJSON 列 unstruct_event_com_acme_checkout_started_1 JSON, unstruct_event_com_acme_product_viewed_1 JSON, -- 时间戳 derived_tstamp TIMESTAMP, load_tstamp TIMESTAMP DEFAULT CURRENT_TIMESTAMP )4.3 示例事件生成逻辑generate_sample_events()的生成规则对应源码事件类型随机从checkout_started与product_viewed中二选一checkout_started的事件event字段记为unstructproduct_viewed记为struct时间分布以「当前时间 - 30 天」为基准每 5 分钟300 秒一条形成真实感的时间序列checkout 事件包含amount、currency、discount_code可为SAVE10/WELCOME20/null以及 items 数组product_id/quantity/priceproduct_viewed 事件包含product_id、categoryelectronics/clothing/books、price用户上下文user_id、user_typefree/premium/enterprise、随机 1~365 天前的registration_date其他字段随机浏览器Chrome/Firefox/Safari、OSWindows/macOS/Linux、国家US/GB/CA、城市等模拟 enrichment 产物。4.4 查询验证# 按事件类型计数 duckdb snowplow_test.duckdb SELECT event, COUNT(*) FROM snowplow.events GROUP BY event # 查看 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 # 查看完整一行验证所有列 duckdb snowplow_test.duckdb SELECT * FROM snowplow.events LIMIT 1也可以在 Python 中直接连接import duckdb conn duckdb.connect(snowplow_test.duckdb) result conn.execute(SELECT * FROM snowplow.events LIMIT 5).fetchall() conn.close()4.5 定制事件数据修改事件类型与字段分布编辑 setup_duckdb.py 中的generate_sample_events()函数可新增事件类型、调整choice([...])的候选值、修改时间间隔或基准时间调整表结构修改create_database()中的CREATE TABLE语句例如增加新的 JSON 上下文列。五、深入 Mock BDP API 服务器5.1 启动方式mock_bdp_server.py 是基于 Flask 的单文件服务参数如下python mock_bdp_server.py --host localhost --port 8081 --debug--host绑定地址默认localhost--port监听端口默认8081--debug开启 Flask debug 模式启动后日志会打印 API 文档地址http://localhost:8081/与健康检查地址http://localhost:8081/health。5.2 端点清单以源码为准方法与路径作用鉴权方式GET /API 文档返回服务名、版本、端点与鉴权说明无GET /health健康检查返回status: healthy与时间戳无GET /organizations/{orgId}/credentials/v3/tokenJWT Token 签发X-Api-Key-IdX-Api-Key请求头GET /organizations/{orgId}/data-structures/v1列出数据结构schema支持filter/vendor/limit/offset参数Authorization: Bearer tokenGET /organizations/{orgId}/data-structures/v1/{hash}按 hash 获取单个数据结构Bearer tokenGET /organizations/{orgId}/data-products/v2列出数据产品tracking plansv2 wrapped 格式Bearer tokenGET /organizations/{orgId}/event-specs/v1列出事件规格wrapped 格式Bearer tokenGET /organizations/{orgId}/tracking-scenarios/v1列出追踪场景wrapped 格式Bearer tokenGET /organizations/{orgId}/users列出组织用户用于 ownership 中initiatorId→ 邮箱解析Bearer token注意两类返回格式差异源码注释中已明确data-structures与users返回直接数组符合 Swagger specdata-products/event-specs/tracking-scenarios返回wrapped 格式data、includes、errors字段。5.3 认证流程Mock 服务器完整复刻了 BDP 的「API Key → JWT Token」两步鉴权# 1) 换取 Token使用 API Key 头 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 # 返回 {accessToken: mock_jwt_token_12345, expiresAt: ...} # 2) 携带 Token 访问数据接口 curl -H Authorization: Bearer mock_jwt_token_12345 \ http://localhost:8081/organizations/test-org-uuid/data-structures/v1源码细节Token 映射表MOCK_TOKENS {test-key-id:test-secret: mock_jwt_token_12345}未匹配的 key 会生成mock_token_api_key_id缺少鉴权头时返回401 {error: Missing authentication headers}Token 的expiresAt为当前时间 1 小时模拟真实 API 的过期行为数据接口缺少/非法的 Bearer 头时返回401 {error: Unauthorized}fixture 文件缺失时返回500 {error: Mock data not found}便于测试异常路径。5.4 扩展 Mock 端点按需编辑mock_bdp_server.py添加新的 Flask 路由仿照现有typed_route装饰器写法它保留了类型注解以兼容 mypy 的disallow_untyped_decorators在fixtures/下创建新的 fixture 文件用load_fixture(filename)加载实现新的 API 行为如分页、过滤、错误注入。六、Fixtures 参考Mock 数据的设计6.1 data_structures_response.json基础版data_structures_response.json 面向基础测试2 个 schemacom.example.page_viewevent与com.example.user_contextentity不含 ownership 数据meta.customData仅含owner/category等简单标记字段定义完整含required、additionalProperties: false、枚举等 JSON Schema 语法。6.2 data_structures_with_ownership.json增强版data_structures_with_ownership.json 面向所有权提取测试3 个核心 schemacheckout_started、product_viewed、user_context另有checkout_completed、session_context等每个 schema 带deployments数组含version、initiator、ts用于 ownership 推断存在多版本 schema如1-0-0、1-1-0以测试版本演进meta.customData承载自定义元数据。Ownership 提取规则源码/文档共同确认createdBy取最旧 deployment 的initiatormodifiedBy取最新 deployment 的initiator字段级 authorship从版本历史推导。6.3 新增一个 Schema Fixture编辑fixtures/data_structures_with_ownership.json向data数组追加条目即可文档中的模板骨架字段均已实际使用{ 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: developercompany.com, ts: 2024-03-01T10:00:00Z } ], data: { self: { vendor: com.acme, name: new_event, version: 1-0-0 }, properties: { field1: { type: string } } } }七、测试场景编排从开发到 CI 的完整矩阵场景 1纯单元测试最快pytest tests/integration/snowplow/test_snowplow.py -v用途开发期快速验证 schema 解析、ownership 提取、golden file 一致性。场景 2DuckDB 仓库测试# 1) 建库 python metadata-ingestion/tests/integration/snowplow/setup/setup_duckdb.py --event-count 500 # 2) 运行带仓库血缘的摄入 datahub ingest -c metadata-ingestion/tests/integration/snowplow/recipes/snowplow_with_duckdb.yml用途验证 warehouse lineage 提取。场景 3完整 Mock BDP 环境端到端# 终端 1启动 Mock 服务器 cd metadata-ingestion/tests/integration/snowplow/setup python mock_bdp_server.py --port 8081 # 终端 2准备 DuckDB 数据 python setup_duckdb.py --event-count 100 # 终端 3执行摄入真实 HTTP 调用 export SNOWPLOW_ORG_IDtest-org-uuid export SNOWPLOW_API_KEY_IDtest-key-id export SNOWPLOW_API_KEYtest-secret datahub ingest -c ../recipes/test_mock_bdp.yml用途端到端验证连接器与 API 交互的完整链路鉴权、抓取、模型解析、落盘。场景 4Iglu 自动发现需要 Docker若需测试 Iglu-only 模式不依赖 BDP从 Iglu Schema Registry 自动发现 schemacd metadata-ingestion/tests/integration/snowplow/setup docker compose -f docker-compose.iglu.yml up -d python setup_iglu.py # 运行 Iglu 自动发现测试 pytest ../test_snowplow.py -k iglu # 或直接执行 Iglu 自动发现 recipe datahub ingest -c ../recipes/test_iglu_autodiscovery.yml # 测试结束关闭 docker compose -f docker-compose.iglu.yml down -vtest_iglu_autodiscovery.yml 的特点是不配置schemas_to_extract依赖iglu_connection.iglu_server_url指向本地 Iglu Serverhttp://localhost:8081通过/api/schemas端点自动发现 schema。八、与测试框架的衔接golden file 机制本地测试环境的「确定性」来自两层保障冻结时间test_snowplow.py通过time_machine.travel(FROZEN_TIME)固定为2024-01-01 00:00:00golden file 比对摄入输出与golden_files/下的期望文件逐一比对并忽略时间戳、systemMetadata.lastObserved、runId等字段见 test_snowplow.py 的ignore_paths列表。若功能变更导致输出有意变化可用以下命令更新 golden 文件pytest tests/integration/snowplow/test_snowplow.py --update-golden-files配合-vv查看详细 diff 定位差异。此外 test_snowplow_performance.py 还覆盖了性能验证并行抓取 deployment、实例级缓存、URN 缓存、1000 个 schema 的大数据集、API 调用次数优化同样可以离线运行。九、故障排查DuckDB 数据库未找到# 确认在 setup 目录下执行 cd metadata-ingestion/tests/integration/snowplow/setup # 重建数据库 python setup_duckdb.pyMock 服务器连接被拒# 检查服务是否存活 curl http://localhost:8081/health # 检查端口占用 lsof -i :8081 # 换端口启动 python mock_bdp_server.py --port 8082Fixture 文件缺失# 确认 fixture 存在应随仓库提交 ls -la metadata-ingestion/tests/integration/snowplow/fixtures/ # 恢复被误删的文件 git checkout metadata-ingestion/tests/integration/snowplow/fixtures/data_structures_with_ownership.json集成测试失败输出有预期内变化时pytest tests/integration/snowplow/test_snowplow.py --update-golden-files输出有意外变化时pytest tests/integration/snowplow/test_snowplow.py -vv查看详细 diff定位是哪一步产生了偏差。Recipe 中 API 前缀注意仓库中部分 recipe 的console_api_url带/api/msc/v1前缀如snowplow_with_duckdb.yml、test_event_specs.yml而 Mock 服务器的路由直接以/organizations/...开头。直接对接 Mock 服务器时请使用 test_mock_bdp.yml 中的写法将console_api_url指向服务器根地址http://localhost:8081。十、环境验证与下一步搭建完成后可参考 SETUP_VERIFICATION.md 按清单逐项核验。推荐的行动顺序跑基础测试pytest metadata-ingestion/tests/integration/snowplow/ -v搭建 DuckDBpython setup_duckdb.py用于仓库测试启动 Mock 服务器python mock_bdp_server.py用于 HTTP 测试定制 fixtures编辑 JSON 匹配你的测试场景参考 OWNERSHIP_TESTING_GUIDE.md跑完整摄入组合 Mock 服务器 DuckDB recipe 完成端到端验证上线前回归最后使用 REAL_BDP_TESTING_GUIDE.md 中描述的 Option A真实 BDP 账号做生产环境前的最终校验并参考 BDP_API_VALIDATION.md 与 SWAGGER_VALIDATION_REPORT.md 核对 API 行为。十一、总结Option B本地测试方案通过JSON Fixture DuckDB Flask Mock 服务器三件套为 Snowplow 连接器提供了自包含、可重复、离线的测试环境其价值在于✅ 快速开发迭代无需等待外部环境5 分钟即可起步✅ CI/CD 友好零外部依赖、确定性输出冻结时间 golden file✅ 离线开发不依赖网络与付费账号✅ 确定性结果mock 数据完全可控便于 diff 与回归✅ 覆盖面完整从纯 Mock 单测、仓库血缘到端到端 HTTP 均可本地复现。当本地环境全部验证通过后再切换到 Option A真实 BDP做最终验收即可在最低成本下获得与生产一致的发布信心。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考