DataHub Mode 数据源接入指南:从 BI 报告到表级血缘的元数据采集实战
2026/9/19 7:24:31 网站建设 项目流程

DataHub Mode 数据源接入指南:从 BI 报告到表级血缘的元数据采集实战

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

Mode 是面向数据分析团队的商业智能(BI)与分析平台。在 DataHub 生态中,mode数据源(连接器)负责将 Mode 中的报告(Report)、图表(Chart)、数据集(Dataset)等 BI 资产,连同所有者信息、表级与列级血缘以及基于状态的有状态删除检测,统一采集进 DataHub 目录。阅读本文后,你将掌握 Mode 连接器的概念映射、前置权限准备、完整配置项语义、典型 Recipe 写法,以及其底层基于 Mode REST API 与 SQL 解析引擎的实现原理,能够独立完成一套可复用的 Mode → DataHub 元数据同步方案。

Overview:DataHub 的 Mode 集成能做什么

Mode 是一个商业智能与分析平台(其官方能力可参考 Mode 官方文档)。DataHub 对 Mode 的集成主要覆盖以下 BI 实体与元数据类型:

  • Dashboards(仪表盘):对应 Mode 中的 Report,包含标题、描述、图表关联、浏览路径、使用统计与可选的嵌入地址(embed URL);
  • Charts(图表):对应 Mode 中挂在 Query 下的可视化,包含图表类型、标题、自定义属性(X/Y 轴、度量、筛选条件等);
  • Datasets(数据集):既包括 Mode 中独立的 Dataset,也包括 Report 内的 Query,二者在 DataHub 中以 Dataset 实体承载;
  • 所有权上下文:将报告/图表的创建者(Creator)映射为 DataHub 的 CorpUser 所有者;
  • 表级与列级血缘:解析 Query 的 SQL,提取上游表与字段依赖,生成粗粒度(table-level)与细粒度(column-level)血缘;
  • 有状态删除检测:基于 DataHub 有状态摄入框架,自动识别 Mode 侧已删除的实体并清理。

从连接器注册表(connector_registry/datahub.json)可以看到,mode连接器的类名为datahub.ingestion.source.mode.ModeSource,支持状态为GA(正式发布),默认启用容器(CONTAINERS)、描述(DESCRIPTIONS)、平台实例(PLATFORM_INSTANCE)、所有权(OWNERSHIP)能力,并默认支持粗粒度与细粒度血缘。

概念映射(Concept Mapping)

原文档指出,Mode 具体的概念映射细节仍有待完善,但 DataHub 有通用的概念映射关系可作参考。下表即这一通用映射:

源概念(Source Concept)DataHub 概念说明
Platform / account / project scopePlatform Instance, Container在平台上下文中组织资产。Mode 集成中对应 Workspace 与 Space(Collection)。
Core technical asset(例如 table / view / topic / file)Dataset主要摄入的技术资产。Mode 集成中对应 Query 与 Mode Dataset。
Schema fields / columnsSchemaField在支持 schema 提取时纳入。Mode 集成通过 SQL 输出列推断生成。
Ownership and collaboration principalsCorpUser, CorpGroup由支持所有权与身份元数据的模块发出。
Dependencies and processing relationshipsLineage edges在支持血缘提取并启用时可用。

结合源码可以给出更具体的映射落地方式(实现见 mode.py):

  • Workspace:作为浏览路径的根(/mode/{workspace}/...),并在_browse_path_space中生成首级 BrowsePathEntry;
  • Space(Collection):映射为 DataHub 容器(Container)。源码中的SpaceKey容器键携带space_token,通过construct_space_container生成容器并打上MODE_COLLECTION子类型(BIContainerSubTypes.MODE_COLLECTION);
  • Report:映射为 Dashboard,子类型MODE_REPORT
  • Query / Mode Dataset:均映射为 Dataset,子类型分别为MODE_QUERYMODE_DATASET
  • Chart:映射为 Chart,子类型MODE_CHART
  • 创作者:映射为urn:li:corpuser:{username}(默认取 username,而非 email)。

前置条件:认证与权限

认证(Authentication)

Mode 连接器使用Basic Auth进行认证,需要一组工作区 API Key:

  • token:创建 Workspace API Key 时的Key ID
  • password:创建 Workspace API Key 时的Secret

从源码看,认证通过requests.SessionHTTPBasicAuth(self.config.token, self.config.password.get_secret_value())完成,同时请求头固定为Content-Type: application/jsonAccept: application/hal+json(Mode API 返回 HAL+JSON 格式)。初始化时连接器会先请求{connect_uri}/api/verify验证连通性与凭据,失败则直接终止并上报Failed to Connect错误。

需要特别注意:Mode 仅支持用户账户认证,不支持服务账户。官方建议为 DataHub 摄入专门创建一个专用用户,避免影响其他成员的日常使用。

权限(Permissions)

Mode 集成要求摄入用户具备以下最小权限(详见 mode_pre.md):

  1. 至少拥有Member角色;
  2. 对每个Connection至少拥有View访问权:可在 "Workspace Settings" → "Manage Connections" 中检查,点击连接 → "Permissions",若默认工作区访问是 "View" 或 "Query" 即可;若为 "Restricted",需要单独给摄入用户授予 View 权限;
  3. 对每个Space(Collection)至少拥有View访问权:以管理员身份进入 "My Collections" 页面,对 Workspace Access 为 "Restricted" 的集合,需在 "Manage Access" 对话框中手动授予摄入用户 "Viewer" 权限;"All Members can View/Edit" 的集合无需手动授权。

注意:若摄入用户拥有Admin权限,则会自动获得所有连接与集合的 View 权限。如果摄入失败,请优先按上述顺序排查凭据、权限、连通性与范围过滤配置,再结合日志中的 source 级错误信息调整。

配置项详解:ModeConfig 全参数参考

Mode 连接器的配置模型为ModeConfig(见 mode.py),继承自有状态摄入配置与数据源血缘公共配置。下表整理了完整参数、默认值与语义:

参数类型默认值说明
connect_uristrhttps://app.mode.comMode 主机地址,结尾多余斜杠会被自动去除。
tokenstr必填Workspace API Key 的 Key ID,用于 Basic Auth。
passwordSecretStr必填Workspace API Key 的 Secret。
workspacestr必填Mode 工作区用户名(URL 中https://app.mode.com/organizations/<workspace-username>部分),与显示名不同,必须全小写
exclude_restrictedboolFalse是否排除 restricted(受限)的集合。源码同时检查restricted字段与default_access_level == "restricted"(因 Mode 侧restricted字段存在已知 bug)。
exclude_personal_collectionsboolTrue是否通过 Mode 服务端过滤器(?filter=custom)排除个人集合。为True时仅拉取共享/自定义集合;为False时拉取全部集合(space_pattern仍会做客户端过滤)。
space_patternAllowDenyPatterndeny: ["^Personal$"]按空间/集合名称做正则过滤,默认排除名为 "Personal" 的集合。例如只想摄入名为 analytics 的空间,可用allow: ['analytics']
report_patternAllowDenyPattern全部允许按报告名称过滤,例如deny: ['slow_report']可排除名为 slow_report 的报告。
owner_username_instead_of_emailboolTrue生成所有者 URN 时使用 username 而非 email。
ingest_embed_urlboolTrue是否为报告生成嵌入地址(embed URL)aspect。
tag_measures_and_dimensionsboolTrue是否为 schema 中的度量(measures)与维度(dimensions)字段打标签。
exclude_archivedboolFalse是否排除已归档(archived)的报告。
max_threadsint1并行 API 请求的最大线程数(1–50)。增大可加速大工作区的摄入,但设置过高可能触发 Mode API 限流(429)。
api_options.retry_backoff_multiplierint/float2指数退避重试的乘数。
api_options.max_retry_intervalint/float60重试等待的最大间隔(秒)。
api_options.max_attemptsint10失败前最大重试次数。
api_options.timeoutint40单次请求等待 Mode REST API 返回数据的超时时间(秒)。
api_options.requests_per_minuteint180所有线程合计的每分钟最大 API 请求数。Mode API 限制约 240 req/min(4 req/s),默认 180 留出余量避免 429。
items_per_pageint30分页请求的每页条数(1–1000),属隐藏配置项。
stateful_ingestionStatefulStaleMetadataRemovalConfigNone有状态摄入与陈旧实体清理配置。

注意:原文档配套示例中出现的default_schema参数已被移除(pydantic_removed_field标注,2025 年 1 月后不再生效),新配置中无需再填写。

API 选项的底层实现

api_options的实现要点(源码 mode.py 与_get_request_json):

  • 请求通过tenacity.Retrying包装,对HTTP 429、HTTP 504 与连接错误wait_exponential(multiplier, max)指数退避、最多max_attempts次重试;
  • 429 时会读取响应头retry-after并休眠等待;504 时短睡 0.1 秒后重试;
  • 所有请求经过RateLimiter(每 60 秒窗口最多requests_per_minute次)限流;
  • 会话级HTTPAdapter还配置了 3 次重试与 10 倍退避因子,连接池大小随max_threads扩容。

快速开始:Recipe 配置示例

官方示例配置见 mode_recipe.yml,完整形态如下(可在其基础上补充sink与可选参数):

source: type: mode config: # 坐标 connect_uri: https://app.mode.com # 凭据(Workspace API Key) token: token password: pass # 选项 workspace: "datahub" owner_username_instead_of_email: false api_options: retry_backoff_multiplier: 2 max_retry_interval: 10 max_attempts: 5 sink: # sink 配置,例如 file 或 datahub-rest type: file config: filename: mode_mces.json

一个更贴近生产、包含过滤与并行加速的示例:

source: type: mode config: connect_uri: https://app.mode.com token: ${MODE_API_TOKEN} password: ${MODE_API_SECRET} workspace: "datahub" exclude_personal_collections: true space_pattern: allow: - "analytics" - "marketing" deny: - "^Personal$" report_pattern: deny: - "slow_report" exclude_archived: true ingest_embed_url: true tag_measures_and_dimensions: true max_threads: 8 api_options: retry_backoff_multiplier: 2 max_retry_interval: 60 max_attempts: 10 timeout: 40 requests_per_minute: 180 stateful_ingestion: enabled: true remove_stale_metadata: true

仓库中还提供了mode连接器的 DHUB 示例(examples/recipes/mode_to_datahub.dhub.yaml),可对照参考。使用datahub ingest子命令加载该 Recipe 即可运行摄入:

datahub ingest -c mode_recipe.yml

摄入完成后,可通过datahubCLI 或 DataHub UI 校验生成的 MCE/MCP 结果。

工作原理:从 REST API 到 DataHub 实体

Mode 连接器基于Mode API v3实现(类注释见 mode.py),核心工作流分三个阶段:先拉取空间并生成容器,再处理报告(Report)及其查询与图表,最后处理独立数据集。

涉及的 Mode API 端点

用途端点
连接验证{connect_uri}/api/verify
空间/集合列表{workspace_uri}/spaces?filter=custom\|all
报告列表(分页){workspace_uri}/spaces/{space_token}/reports?filter=all&per_page=N&page=N
数据集列表(分页){workspace_uri}/spaces/{space_token}/datasets?filter=all&per_page=N&page=N
报告查询列表{workspace_uri}/reports/{report_token}/queries
查询图表列表{workspace_uri}/reports/{report_token}/queries/{query_token}/charts
数据源(Connection)列表{workspace_uri}/data_sources
定义(Definitions)列表{workspace_uri}/definitions

其中{workspace_uri}{connect_uri}/api/{workspace}。源码注释指出查询与图表两个端点不支持分页,直接整体拉取;空间、报告、数据集则通过_get_paged_request_jsonitems_per_page分页遍历。

实体构建流程

  • 容器construct_space_container为每个 Space 生成MODE_COLLECTION容器,并写入以工作区为根的 BrowsePathsV2;
  • 报告 → Dashboardconstruct_dashboard使用报告的名称、描述、创建/修改时间戳(created_at/last_saved_at,后者为空时回退edited_at)、图表 URN 列表与引用数据集构建DashboardInfoClass,同时生成浏览路径、MODE_REPORT子类型、所有权(DATAOWNER)、使用统计(view_count)与可选的 embed 链接;
  • 查询/数据集 → Datasetconstruct_query_or_dataset生成DatasetPropertiesClass(外部 URL、自定义属性如idcreated_atdata_source_idchart_count等)、ViewPropertiesClass(原始 SQL 视图逻辑)、SubTypesClass与 BrowsePathsV2,并解析 SQL 产出 Schema 元数据与血缘;
  • 图表 → Chartconstruct_chart_from_api_data根据 Mode 图表的chartType/selectedChart映射 DataHub 图表类型(见_get_chart_type映射表),并抽取encoding中的 X/Y/度量/筛选公式写入自定义属性,同时生成ChartQueryClass(原始 SQL)、InputFields(字段级输入关系)与所有权。

图表类型映射表(_get_chart_type)值得一提:Mode 的table/pivotTable/vegasPivotTable→ DataHubTABLEbar/stackedBar/stackedBar100/hStackedBar/hStackedBar100/hBarBARlineLINEarea/totalAreaAREApie/donutPIEscatterSCATTERbigNumber/bigValueTEXThistogramHISTOGRAM;未识别的类型置为None并在报告中记录 warning。

SQL 血缘提取:三层预处理

表级与列级血缘是 Mode 集成的核心能力,其提取链路在construct_query_or_datasetget_upstream_lineage_for_parsed_sql中实现,包含三层 SQL 预处理:

  1. 定义展开(Definitions):Mode 查询中常使用{{ @definition as alias }}语法引用公共定义。_replace_definitions会从/definitions端点拉取定义源,递归替换(最大深度 10,检测循环引用),还原成完整 SQL;
  2. 表单模板渲染(Liquid):Mode 查询支持{% form %}...{% endform %}块与 Liquid 模板变量。normalize_mode_query解析表单参数的默认值并渲染模板,若渲染失败则退回去掉表单块后的 SQL,保证后续解析不因模板语法报错;
  3. SQL 解析(sqlglot):使用sqlglot_lineage对多语句 SQL 取最后一个有效语句进行解析,结合从 DataHub Graph 缓存构建的 schema resolver,产出SqlParsingResult。随后infer_output_schema推断查询输出列生成SchemaMetadataClass,并生成粗粒度(UpstreamLineageClass,类型TRANSFORMED)与细粒度(FineGrainedLineageClass,FIELD → FIELD_SET)血缘。

此外还有两处平台相关的特殊处理:

  • JDBC 适配器映射_get_datahub_friendly_platform将 Mode 数据源(Connection)的 JDBC 前缀映射为 DataHub 平台名,例如jdbc:athena → athenajdbc:bigquery → bigqueryjdbc:postgresql → postgresjdbc:snowflake → snowflakejdbc:redshift → redshift等;未识别的前缀使用原始名称并记录 warning;
  • Snowflake Warehouse 识别:当上游平台为 Snowflake 时,通过正则从 SQL 中提取use warehouse <name>;语句,将其作为默认数据库名参与血缘解析。

度量/维度标签启发式

启用tag_measures_and_dimensions后,set_field_tags会为 schema 字段打上DimensionMeasure标签。其启发式规则为:字段类型为数值(NumberType)且字段名不以_number结尾、也不是 id 类字段(匹配(^id[_\d]?)|([_\d?]id$))时,标记为Measure(度量);否则标记为Dimension(维度)。源码注释明确说明该启发式"可能不准确",因为 Mode 侧目前无法明确区分度量与维度。

有状态摄入与删除检测

ModeSource继承StatefulIngestionSourceBase,配置stateful_ingestion.remove_stale_metadata: true后启用陈旧实体清理,可自动检测并移除 Mode 侧已删除的报告、查询、图表等实体(对应 README 中提到的 stateful deletion detection)。该能力在连接器注册表中也标记为 "Enabled by default via stateful ingestion"。

性能与并发设计

针对大型工作区,连接器提供了以下性能机制(均有源码与测试佐证):

  • 全局线程池max_threads > 1时,报告与数据集处理分别通过ThreadedIteratorExecutor在跨空间的全局线程池中并行执行(见get_workunits_internal);
  • 限流与重试RateLimiter保证全局每分钟请求数不超阈值;429/504/连接错误自动指数退避重试;
  • 缓存:数据源列表、定义列表、创建者信息(serialized_lru_cache,maxsize 5000)均有缓存,且只缓存成功结果,避免瞬时失败被永久缓存;
  • 报告级容错:单个报告处理失败或超时仅记录 warning/failure,不会中断整个摄入管道;报告级错误处理本身还有二次异常保护,防止错误上报自身抛错导致线程池整体中止;
  • 跳过空图表 API 调用:当查询的chart_count == 0时直接跳过图表接口调用,减少无效请求;
  • 内存观测:摄入结束后在报告中记录进程 RSS 内存占用(process_memory_used_mb),便于监控大工作区摄入的资源消耗。

并发执行的安全细节也有专门处理:报告对象(ModeSourceReport)内部使用threading.Lock保护结构化日志与计数器,PerfTimer被明确标注为非线程安全、只能在主线程使用。

能力边界与限制

根据 mode_post.md 与源码实现,Mode 连接器的能力与限制总结如下:

能力说明
报告(Report)通过 Mode 报告 API 获取标题、描述、所有权与图表关联。
图表(Chart)通过reports/{report}/queries/{query}/charts获取,含图表类型、标题与构建 DataHub Chart 实体所需的元数据。
表信息(Table)报告查询的表结果元数据用于识别上游数据集上下文与查询关系。
透视表信息(Pivot Table)可用时提取透视结果元数据,改善透视分析的图表/数据集关系覆盖。

需要说明的限制与注意事项:

  • 连接器仅实测过 PostgreSQL 数据源,其他数据库理论上可用但未经验证(源码类注释明确说明);
  • 血缘解析依赖 SQL 可解析性:定义循环引用、模板渲染失败、数据源连接被删除(data_source_id找不到)等场景都会导致血缘被跳过并在报告中记录告警;
  • 透视表/图表字段关联基于公式中[...]引用的列名进行大小写不敏感匹配,匹配不到真实 schema 字段时会跳过;
  • 图表类型与度量/维度标签均存在"未知类型置空"与"启发式判定"的降级策略;
  • 摄入行为受 Mode API、权限与平台暴露的元数据约束,mode_post.md建议以能力表为权威依据。

验证与测试

仓库为 Mode 连接器提供了完整的集成测试,可作为理解与验证行为的参考:

  • tests/integration/mode/test_mode.py:通过 Mock 响应模拟 Mode API(/api/verify、spaces、reports、queries、charts、data_sources、definitions 等端点),验证完整摄入管线的 MCE 输出;
  • tests/integration/mode/test_mode_threading.py:验证max_threads并发处理场景;
  • tests/integration/mode/mode_mces_golden.json:摄入结果的 golden 期望文件,展示了报告、图表、数据集实体的最终 MCE 形态。

测试中的JSON_RESPONSE_MAP清晰地还原了连接器实际调用的 API 路径序列(如https://app.mode.com/api/acryl/spaces/157933cc1168/reports.../reports/9d2da37fa91e/queries.../queries/6e26a9f3d4e2/charts等),与上文端点表一一对应。

结语

Mode 连接器是 DataHub BI 生态中较完整的开源实现:它既覆盖了报告、图表、数据集等 BI 资产的基础元数据,也通过定义展开、Liquid 模板渲染与 sqlglot 解析实现了可用的表级与列级血缘,并辅以所有权、容器、使用统计与有状态删除检测等增强能力。部署时建议遵循"专用摄入用户 + 最小权限 + 空间/报告过滤 + 限流参数"的最佳实践组合,即可稳定地将 Mode 中的分析资产接入 DataHub 统一目录。

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

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

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

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

立即咨询