ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

DataHub SnapLogic 集成指南:流式与集成实体元数据及表列级血缘接入实战

DataHub SnapLogic 集成指南:流式与集成实体元数据及表列级血缘接入实战 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-levellineage元数据类型包括数据集含 schema 字段、PipelineData Flow、SnapData Job。该插件在仓库中的实现位于metadata-ingestion/src/datahub/ingestion/source/snaplogic/目录核心入口为SnaplogicSource。插件注册信息见 metadata-ingestion/src/datahub/ingestion/autogenerated/connector_registry/datahub.json其中 platform_id 为snaplogicplatform_name 为SnapLogic支持状态为ALPHA。1.1 能力矩阵Capabilities从 snaplogic.py 的装饰器声明与注册表可以确认当前能力边界能力支持情况说明LINEAGE_COARSE粗粒度血缘默认启用Pipeline/Snap 与数据集之间的血缘LINEAGE_FINE细粒度血缘默认启用字段级列级血缘PLATFORM_INSTANCE平台实例不支持SnapLogic 不支持 platform instancesDELETION_DETECTION删除检测暂不支持不检测源端删除实体二、核心概念映射SnapLogic 实体 → DataHub 实体这是本集成最重要的设计骨架。原文档给出了如下映射关系直接决定了产出的 DataHub 元数据形态Source Concept源概念DataHub ConceptDataHub 概念备注Snap-packData PlatformSnap-packs 映射为 Data Platform既可以是直接映射如 Snowflake也可以根据连接信息动态推导如 JDBC URL。Table/DatasetDataset可能有所不同取决于 Snap 类型对 SQL 数据库是表table对 Kafka 则是主题topic。SnapData Job每个 Snap 映射为一个 Data Job。PipelineData Flow每个 Pipeline 映射为一个 Data Flow。2.1 从源码看映射如何落地上述映射在代码中由SnapLogicParser与SnaplogicSource共同落实Pipeline → Data Flowcreate_pipeline_mcp()snaplogic.py用make_data_flow_urn(orchestratornamespace, flow_idpipeline_snode_id, clusterPROD)生成 Data Flow URN并写入DataFlowInfoClassname、externalUrl。Snap → Data Jobcreate_task_mcp()snaplogic.py用make_data_job_urn(orchestratornamespace, flow_idpipeline_snode_id, job_idtask_id, clusterPROD)生成 Data Job URNDataJobInfoClass.type SNAPLOGIC_SNAP。Table/Dataset → Datasetcreate_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.pysnaplogic: set()与 pyproject.toml入口点entry point注册于 setup.py 与 pyproject.tomlsnaplogic 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: examplesnaplogic.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字段说明如下配置项类型必填默认值说明usernamestring是—SnapLogic 用户名passwordSecretStr是—SnapLogic 密码Secret 类型落盘时被脱敏base_urlstring否https://elastic.snaplogic.comSnapLogic 实例地址用于调用其 APIorg_namestring是—SnapLogic 实例中的组织Organization名称namespace_mappingdict否{}namespace 到 platform instance 的映射case_insensitive_namespaceslist否[]需要按大小写不敏感处理的 namespace 列表create_non_snaplogic_datasetsbool否False是否为非 SnapLogic 平台的数据集数据库、S3 等创建 Dataset 实体stateful_ingestionobject否None有状态摄取配置见下文platformstring否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.pycreate_dataset_mcp中的跳过逻辑。stateful_ingestion支持两个子项enabled开启有状态摄取remove_stale_metadata是否清理过期元数据示例中为False且当前插件尚未实现删除检测能力建议保持关闭。SnaplogicConfig同时混入了StatefulLineageConfigMixin与StatefulUsageConfigMixinsnaplogic_config.py因此还支持有状态血缘相关的start_time、end_time、enable_stateful_lineage_ingestion等配置用于限定血缘查询的时间窗口与跳过冗余运行。六、血缘提取原理OpenLineage API 拉取与分页血缘提取由SnaplogicLineageExtractorsnaplogic_lineage_extractor.py完成其核心是get_lineages()L31-L87请求端点GET {base_url}/api/1/rest/public/catalog/{org_name}/lineage查询参数formatOPENLINEAGE返回 OpenLineage 格式的 RunEventstart_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_snode685013b9da1804dd3b4037e8, eventType: COMPLETE, run: { facets: { parent: { _producer: ...?#pipe_snode685013b9da1804dd3b4037e8, 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_snodeid标识所属 Pipeline对应 Data Flowinputs/outputs是数据集Datasetfacets.schema.fields携带字段与类型outputs[].facets.columnLineage.fields携带列级血缘每个输出字段列出其inputFields来源数据集 来源字段。6.2 单条记录的加工链路SnaplogicSource.get_workunits_internal()snaplogic.py逐条拉取血缘记录每 20 条输出一次进度日志单条记录经_process_lineage_record()L132-L181处理从producer中解析pipe_snodePipeline ID缺失则跳过由SnapLogicParser抽取数据集含 INPUT/OUTPUT 类型标注、Pipeline、Snap 任务、列映射依次产出 PipelineData FlowMCP、Dataset MCP、TaskData JobMCP含血缘。七、列级血缘与类型映射细节7.1 列级血缘Fine-Grained Lineagecreate_task_mcp()snaplogic.py在DataJobInputOutputClass中填写inputDatasets/outputDatasets粗粒度血缘填写inputDatasetFields/outputDatasetFields数据集字段全集为每个ColumnMapping生成FineGrainedLineageClassupstream/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、varcharStringTypeClassnumber、long、float、double、intNumberTypeClassbooleanBooleanTypeClass其他兜底StringTypeClass映射后的 schema 与原生类型nativeDataType一起写入SchemaMetadataClasssnaplogic.py。7.3 外部 URL 与跳转Pipeline 与 Data Job 的externalUrl均为{base_url}/sl/designer.html?v21818#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.jsontest_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),仅供参考
返回列表