DataHub SnapLogic 集成指南:流式与集成实体元数据及表列级血缘接入实战
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
本文围绕 DataHub 官方 SnapLogic 数据源(Source)插件展开,介绍如何将 SnapLogic(流式/集成平台)中的 topics、connectors、pipelines、jobs 等流式与集成实体及其表级、列级血缘同步到 DataHub。读者读完本文,将掌握 SnapLogic 与 DataHub 的实体概念映射、完整 recipe 配置、底层血缘提取原理(OpenLineage API 拉取、分页与时间窗口),并能基于仓库源码与集成测试快速上手、排障与二次开发。
一、插件概览:SnapLogic 元数据能同步什么
Snaplogic 是一个流式(streaming)或集成(integration)平台,官方资料可参考其产品文档。DataHub 为其提供的集成插件snaplogic覆盖以下内容(见 metadata-ingestion/docs/sources/snaplogic/README.md):
- 流式/集成实体:topics、connectors、pipelines、jobs;
- 血缘:同时捕获表级(table-level)与列级(column-level)lineage;
- 元数据类型:包括数据集(含 schema 字段)、Pipeline(Data Flow)、Snap(Data Job)。
该插件在仓库中的实现位于metadata-ingestion/src/datahub/ingestion/source/snaplogic/目录,核心入口为SnaplogicSource。插件注册信息见 metadata-ingestion/src/datahub/ingestion/autogenerated/connector_registry/datahub.json,其中 platform_id 为snaplogic,platform_name 为SnapLogic,支持状态为ALPHA。
1.1 能力矩阵(Capabilities)
从 snaplogic.py 的装饰器声明与注册表可以确认当前能力边界:
| 能力 | 支持情况 | 说明 |
|---|---|---|
| LINEAGE_COARSE(粗粒度血缘) | 默认启用 | Pipeline/Snap 与数据集之间的血缘 |
| LINEAGE_FINE(细粒度血缘) | 默认启用 | 字段级(列级)血缘 |
| PLATFORM_INSTANCE(平台实例) | 不支持 | SnapLogic 不支持 platform instances |
| DELETION_DETECTION(删除检测) | 暂不支持 | 不检测源端删除实体 |
二、核心概念映射:SnapLogic 实体 → DataHub 实体
这是本集成最重要的设计骨架。原文档给出了如下映射关系,直接决定了产出的 DataHub 元数据形态:
| Source Concept(源概念) | DataHub Concept(DataHub 概念) | 备注 |
|---|---|---|
| Snap-pack | Data Platform | Snap-packs 映射为 Data Platform,既可以是直接映射(如 Snowflake),也可以根据连接信息动态推导(如 JDBC URL)。 |
| Table/Dataset | Dataset | 可能有所不同,取决于 Snap 类型:对 SQL 数据库是表(table),对 Kafka 则是主题(topic)。 |
| Snap | Data Job | 每个 Snap 映射为一个 Data Job。 |
| Pipeline | Data Flow | 每个 Pipeline 映射为一个 Data Flow。 |
2.1 从源码看映射如何落地
上述映射在代码中由SnapLogicParser与SnaplogicSource共同落实:
- Pipeline → Data Flow:
create_pipeline_mcp()(snaplogic.py)用make_data_flow_urn(orchestrator=namespace, flow_id=pipeline_snode_id, cluster="PROD")生成 Data Flow URN,并写入DataFlowInfoClass(name、externalUrl)。 - Snap → Data Job:
create_task_mcp()(snaplogic.py)用make_data_job_urn(orchestrator=namespace, flow_id=pipeline_snode_id, job_id=task_id, cluster="PROD")生成 Data Job URN,DataJobInfoClass.type = "SNAPLOGIC_SNAP"。 - Table/Dataset → Dataset:
create_dataset_mcp()(snaplogic.py)通过make_dataset_urn_with_platform_instance()生成 Dataset URN,同时产出DatasetPropertiesClass与SchemaMetadataClass(字段级 schema)。 - Snap-pack → Data Platform:由 snaplogic_parser.py 的
_parse_platform()动态解析:取 namespace 中://前的协议部分并转小写作为平台名(例如sqlserver://...→sqlserver,并进一步通过platform_mapping映射为 DataHub 的mssql)。
值得注意:namespace 即来自 SnapLogic Lineage API 中的 OpenLineage 格式namespace字段。例如集成测试样例 snaplogic_simple_response.json 中,输入数据集 namespace 为sqlserver://snaplogic-test.database.windows.net:1433,会被解析为mssql平台的表snaplogic-test.tonyschema.accounts;而输出数据集 namespace 为SnapLogic,属于 SnapLogic 平台内部的虚拟表。
三、前置条件(Prerequisites)
模块说明见 snaplogic_pre.md,运行摄取前需要满足:
- 确保能访问 SnapLogic 源实例的网络连通性;
- 拥有有效的 SnapLogic 认证凭据;
- 具备读取本模块所需元数据 API 的权限;
- 必须有访问 SnapLogic Lineage API 的有效凭据——因为血缘数据全部来自该 API。
四、安装与启用
snaplogic作为acryl-datahub的插件模块注册:
- 可选依赖(extra)定义于 setup.py(
"snaplogic": set())与 pyproject.toml; - 入口点(entry point)注册于 setup.py 与 pyproject.toml:
snaplogic = datahub.ingestion.source.snaplogic.snaplogic:SnaplogicSource。
因此安装方式为:
pip install 'acryl-datahub[snaplogic]'安装后即可在 recipe 中声明type: snaplogic使用。
五、完整配置详解:Recipe 与参数说明
官方示例 recipe 位于 snaplogic_recipe.yml,内容如下(这也是仓库中唯一一份可复制的完整配置样例):
pipeline_name: "snaplogic_incremental_ingestion" source: type: snaplogic config: username: example@snaplogic.com password: password base_url: https://elastic.snaplogic.com org_name: "ExampleOrg" namespace_mapping: snowflake://snaplogic: snaplogic case_insensitive_namespaces: - snowflake://snaplogic stateful_ingestion: enabled: True remove_stale_metadata: False5.1 配置字段对照表
配置模型定义于 snaplogic_config.py,字段说明如下:
| 配置项 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
username | string | 是 | — | SnapLogic 用户名 |
password | SecretStr | 是 | — | SnapLogic 密码(Secret 类型,落盘时被脱敏) |
base_url | string | 否 | https://elastic.snaplogic.com | SnapLogic 实例地址,用于调用其 API |
org_name | string | 是 | — | SnapLogic 实例中的组织(Organization)名称 |
namespace_mapping | dict | 否 | {} | namespace 到 platform instance 的映射 |
case_insensitive_namespaces | list | 否 | [] | 需要按大小写不敏感处理的 namespace 列表 |
create_non_snaplogic_datasets | bool | 否 | False | 是否为非 SnapLogic 平台的数据集(数据库、S3 等)创建 Dataset 实体 |
stateful_ingestion | object | 否 | None | 有状态摄取配置(见下文) |
platform | string | 否 | SnapLogic | 平台名(内部固定值) |
5.2 关键参数的底层影响
namespace_mapping:在SnapLogicParser构造时传入(snaplogic.py),用于在_create_dataset_info()(snaplogic_parser.py)中为数据集附加platform_instance。典型用途:把形如snowflake://snaplogic的 namespace 归一到某个 DataHub 平台实例。case_insensitive_namespaces:当某个 namespace 出现在该列表中,数据集名与字段名会被统一转为小写(snaplogic_parser.py、L137-L166),避免大小写差异导致重复实体。create_non_snaplogic_datasets:控制是否把血缘中出现的非 SnapLogic 数据集(如外部数据库表、S3)也落成 Dataset。默认False时只建 SnapLogic 侧的数据集;设为True后,会为外部平台数据集创建实体,但前提是该数据集尚未存在于 DataHub(见 snaplogic.py,create_dataset_mcp中的跳过逻辑)。stateful_ingestion:支持两个子项:enabled:开启有状态摄取;remove_stale_metadata:是否清理过期元数据(示例中为False,且当前插件尚未实现删除检测能力,建议保持关闭)。
SnaplogicConfig同时混入了StatefulLineageConfigMixin与StatefulUsageConfigMixin(snaplogic_config.py),因此还支持有状态血缘相关的start_time、end_time、enable_stateful_lineage_ingestion等配置,用于限定血缘查询的时间窗口与跳过冗余运行。
六、血缘提取原理:OpenLineage API 拉取与分页
血缘提取由SnaplogicLineageExtractor(snaplogic_lineage_extractor.py)完成,其核心是get_lineages()(L31-L87):
- 请求端点:
GET {base_url}/api/1/rest/public/catalog/{org_name}/lineage; - 查询参数:
format=OPENLINEAGE(返回 OpenLineage 格式的 RunEvent);start_ts/end_ts(毫秒时间戳,限定血缘时间窗口);page(从 0 开始的分页游标);
- 认证:HTTP Basic Auth(用户名 + 密码),并携带
User-Agent: datahub-connector/1.0; - 分页策略:若当前页返回记录数
>= 20,则继续请求下一页,直到不足一页为止(L70-L74); - 逐条产出:以生成器方式逐条
yield血缘记录,交给上层SnaplogicSource处理,避免一次性加载全量数据。
时间窗口的确定见_get_time_window()(L89-L95):开启有状态血缘时,由RedundantLineageRunSkipHandler.suggest_run_time_window()基于上次检查点建议窗口,否则使用配置中的start_time/end_time。这解释了示例 recipe 中pipeline_name取snaplogic_incremental_ingestion的用意——结合有状态摄取实现增量拉取。
6.1 血缘记录的 OpenLineage 结构
集成测试的 mock 数据 snaplogic_simple_response.json 展示了单条记录的关键结构:
{ "producer": "https://tahoe.elastic.snaplogicdev.com/sl/designer.html?#pipe_snode=685013b9da1804dd3b4037e8", "eventType": "COMPLETE", "run": { "facets": { "parent": { "_producer": "...?#pipe_snode=685013b9da1804dd3b4037e8", "job": { "namespace": "SnapLogic", "name": "Datahub Demo 3" } } } }, "job": { "namespace": "SnapLogic", "name": "Datahub Demo 3:Azure Synapse SQL - Select:2faf1220-..." }, "inputs": [{ "namespace": "sqlserver://snaplogic-test.database.windows.net:1433", "name": "snaplogic-test.tonyschema.accounts", "facets": { "schema": { "fields": [ { "name": "Id", "type": "VARCHAR" }, { "name": "Name", "type": "VARCHAR" } ] } } }], "outputs": [{ "namespace": "SnapLogic", "name": "Virtual_DB.Virtual_Schema.Azure Synapse SQL - Select:2faf1", "facets": { "schema": { "fields": [ { "name": "Id", "type": "VARCHAR" }, { "name": "Name", "type": "VARCHAR" } ] }, "columnLineage": { "fields": { "Id": { "inputFields": [ { "field": "Id", "name": "snaplogic-test.tonyschema.accounts", "namespace": "sqlserver://..." } ] } } } } }] }可以看出:
- 顶层的
job是 Snap 任务(对应 Data Job),其producer中的#pipe_snode=<id>标识所属 Pipeline(对应 Data Flow); inputs/outputs是数据集(Dataset),facets.schema.fields携带字段与类型;outputs[].facets.columnLineage.fields携带列级血缘:每个输出字段列出其inputFields(来源数据集 + 来源字段)。
6.2 单条记录的加工链路
SnaplogicSource.get_workunits_internal()(snaplogic.py)逐条拉取血缘记录,每 20 条输出一次进度日志;单条记录经_process_lineage_record()(L132-L181)处理:
- 从
producer中解析pipe_snode(Pipeline ID),缺失则跳过; - 由
SnapLogicParser抽取数据集(含 INPUT/OUTPUT 类型标注)、Pipeline、Snap 任务、列映射; - 依次产出 Pipeline(Data Flow)MCP、Dataset MCP、Task(Data Job)MCP(含血缘)。
七、列级血缘与类型映射细节
7.1 列级血缘(Fine-Grained Lineage)
create_task_mcp()(snaplogic.py)在DataJobInputOutputClass中:
- 填写
inputDatasets/outputDatasets(粗粒度血缘); - 填写
inputDatasetFields/outputDatasetFields(数据集字段全集); - 为每个
ColumnMapping生成FineGrainedLineageClass,upstream/downstream 均为FIELD_SET类型,指向具体的make_schema_field_urn()字段 URN。
ColumnMapping由extract_columns_mapping_from_lineage()(snaplogic_parser.py)从columnLineagefacet 中解析:遍历每个输出字段,逐条关联其inputFields。
7.2 数据类型映射
SnaplogicUtils.get_datahub_type()(snaplogic_utils.py)将 SnapLogic/数据库字符串类型映射为 DataHubSchemaFieldDataTypeClass:
| 源类型(小写后) | DataHub 类型 |
|---|---|
string、varchar | StringTypeClass |
number、long、float、double、int | NumberTypeClass |
boolean | BooleanTypeClass |
| 其他(兜底) | StringTypeClass |
映射后的 schema 与原生类型(nativeDataType)一起写入SchemaMetadataClass(snaplogic.py)。
7.3 外部 URL 与跳转
Pipeline 与 Data Job 的externalUrl均为{base_url}/sl/designer.html?v=21818#pipe_snode={pipeline_snode_id}(snaplogic.py、L332),可在 DataHub 界面直接跳回 SnapLogic Designer 定位对应 Pipeline。
八、测试验证:如何确认集成行为
仓库为 SnapLogic 插件提供了完整的集成测试与黄金文件,位于metadata-ingestion/tests/integration/snaplogic/:
- test_snaplogic.py:用
requests_mock拦截https://elastic.snaplogic.com/api/1/rest/public/catalog/TEST_ORG/lineage,通过mce_helpers.check_golden_file()对比生成的 MCE 与黄金文件,覆盖:- 默认配置下的摄取结果(snaplogic_base_golden.json);
- 开启
create_non_snaplogic_datasets后的结果(snaplogic_create_non_snaplogic_datasets_golden.json);
- test_snaplogic_utils.py:验证类型映射;
- test_snaplogic_lineage_extractor.py:验证血缘提取与解析;
- snaplogic_base_recipe.yml 与 snaplogic_simple_response.json:作为测试输入样例。
九、运行方式与注意事项
9.1 运行摄取
配置好 recipe 后,使用 DataHub CLI 执行:
datahub ingest -c snaplogic_recipe.yml9.2 注意事项
- 权限:必须拥有访问 SnapLogic Lineage API 的凭据与读取权限,否则
get_lineages()请求将失败(raise_for_status()会抛出异常,并在报告中记录 "Error fetching lineage data")。 - 支持状态为 ALPHA:该插件当前为 ALPHA 状态,未实现删除检测(DELETION_DETECTION 不支持)、不支持平台实例(PLATFORM_INSTANCE 不支持),生产环境接入前需评估。
- 非 SnapLogic 数据集默认不建实体:血缘中外部平台(数据库、S3、Kafka 等)的数据集默认仅作为血缘节点引用,如需在 DataHub 中创建对应 Dataset,须显式开启
create_non_snaplogic_datasets。 - 大小写敏感问题:对于
case_insensitive_namespaces中列出的 namespace,数据集名与字段名会统一小写,注意与目标平台实体命名保持一致,避免产生重复实体。 - 平台名归一:
sqlserver://等 namespace 会被归一为 DataHub 平台标识(如mssql),这决定了血缘连到哪个平台的数据集上。
十、延伸阅读
- 实体模型文档:Data Platform、Dataset、Data Job、Data Flow;
- 插件官方说明:snaplogic_pre.md、snaplogic_recipe.yml;
- 核心源码:snaplogic.py、snaplogic_config.py、snaplogic_lineage_extractor.py、snaplogic_parser.py、snaplogic_utils.py;
- 集成测试:test_snaplogic.py、test_data/snaplogic_simple_response.json;
- 平台注册信息:connector_registry/datahub.json。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考