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 scope | Platform Instance, Container | 在平台上下文中组织资产。Mode 集成中对应 Workspace 与 Space(Collection)。 |
| Core technical asset(例如 table / view / topic / file) | Dataset | 主要摄入的技术资产。Mode 集成中对应 Query 与 Mode Dataset。 |
| Schema fields / columns | SchemaField | 在支持 schema 提取时纳入。Mode 集成通过 SQL 输出列推断生成。 |
| Ownership and collaboration principals | CorpUser, CorpGroup | 由支持所有权与身份元数据的模块发出。 |
| Dependencies and processing relationships | Lineage 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_QUERY与MODE_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.Session的HTTPBasicAuth(self.config.token, self.config.password.get_secret_value())完成,同时请求头固定为Content-Type: application/json与Accept: application/hal+json(Mode API 返回 HAL+JSON 格式)。初始化时连接器会先请求{connect_uri}/api/verify验证连通性与凭据,失败则直接终止并上报Failed to Connect错误。
需要特别注意:Mode 仅支持用户账户认证,不支持服务账户。官方建议为 DataHub 摄入专门创建一个专用用户,避免影响其他成员的日常使用。
权限(Permissions)
Mode 集成要求摄入用户具备以下最小权限(详见 mode_pre.md):
- 至少拥有Member角色;
- 对每个Connection至少拥有View访问权:可在 "Workspace Settings" → "Manage Connections" 中检查,点击连接 → "Permissions",若默认工作区访问是 "View" 或 "Query" 即可;若为 "Restricted",需要单独给摄入用户授予 View 权限;
- 对每个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_uri | str | https://app.mode.com | Mode 主机地址,结尾多余斜杠会被自动去除。 |
token | str | 必填 | Workspace API Key 的 Key ID,用于 Basic Auth。 |
password | SecretStr | 必填 | Workspace API Key 的 Secret。 |
workspace | str | 必填 | Mode 工作区用户名(URL 中https://app.mode.com/organizations/<workspace-username>部分),与显示名不同,必须全小写。 |
exclude_restricted | bool | False | 是否排除 restricted(受限)的集合。源码同时检查restricted字段与default_access_level == "restricted"(因 Mode 侧restricted字段存在已知 bug)。 |
exclude_personal_collections | bool | True | 是否通过 Mode 服务端过滤器(?filter=custom)排除个人集合。为True时仅拉取共享/自定义集合;为False时拉取全部集合(space_pattern仍会做客户端过滤)。 |
space_pattern | AllowDenyPattern | deny: ["^Personal$"] | 按空间/集合名称做正则过滤,默认排除名为 "Personal" 的集合。例如只想摄入名为 analytics 的空间,可用allow: ['analytics']。 |
report_pattern | AllowDenyPattern | 全部允许 | 按报告名称过滤,例如deny: ['slow_report']可排除名为 slow_report 的报告。 |
owner_username_instead_of_email | bool | True | 生成所有者 URN 时使用 username 而非 email。 |
ingest_embed_url | bool | True | 是否为报告生成嵌入地址(embed URL)aspect。 |
tag_measures_and_dimensions | bool | True | 是否为 schema 中的度量(measures)与维度(dimensions)字段打标签。 |
exclude_archived | bool | False | 是否排除已归档(archived)的报告。 |
max_threads | int | 1 | 并行 API 请求的最大线程数(1–50)。增大可加速大工作区的摄入,但设置过高可能触发 Mode API 限流(429)。 |
api_options.retry_backoff_multiplier | int/float | 2 | 指数退避重试的乘数。 |
api_options.max_retry_interval | int/float | 60 | 重试等待的最大间隔(秒)。 |
api_options.max_attempts | int | 10 | 失败前最大重试次数。 |
api_options.timeout | int | 40 | 单次请求等待 Mode REST API 返回数据的超时时间(秒)。 |
api_options.requests_per_minute | int | 180 | 所有线程合计的每分钟最大 API 请求数。Mode API 限制约 240 req/min(4 req/s),默认 180 留出余量避免 429。 |
items_per_page | int | 30 | 分页请求的每页条数(1–1000),属隐藏配置项。 |
stateful_ingestion | StatefulStaleMetadataRemovalConfig | None | 有状态摄入与陈旧实体清理配置。 |
注意:原文档配套示例中出现的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_json按items_per_page分页遍历。
实体构建流程
- 容器:
construct_space_container为每个 Space 生成MODE_COLLECTION容器,并写入以工作区为根的 BrowsePathsV2; - 报告 → Dashboard:
construct_dashboard使用报告的名称、描述、创建/修改时间戳(created_at/last_saved_at,后者为空时回退edited_at)、图表 URN 列表与引用数据集构建DashboardInfoClass,同时生成浏览路径、MODE_REPORT子类型、所有权(DATAOWNER)、使用统计(view_count)与可选的 embed 链接; - 查询/数据集 → Dataset:
construct_query_or_dataset生成DatasetPropertiesClass(外部 URL、自定义属性如id、created_at、data_source_id、chart_count等)、ViewPropertiesClass(原始 SQL 视图逻辑)、SubTypesClass与 BrowsePathsV2,并解析 SQL 产出 Schema 元数据与血缘; - 图表 → Chart:
construct_chart_from_api_data根据 Mode 图表的chartType/selectedChart映射 DataHub 图表类型(见_get_chart_type映射表),并抽取encoding中的 X/Y/度量/筛选公式写入自定义属性,同时生成ChartQueryClass(原始 SQL)、InputFields(字段级输入关系)与所有权。
图表类型映射表(_get_chart_type)值得一提:Mode 的table/pivotTable/vegasPivotTable→ DataHubTABLE;bar/stackedBar/stackedBar100/hStackedBar/hStackedBar100/hBar→BAR;line→LINE;area/totalArea→AREA;pie/donut→PIE;scatter→SCATTER;bigNumber/bigValue→TEXT;histogram→HISTOGRAM;未识别的类型置为None并在报告中记录 warning。
SQL 血缘提取:三层预处理
表级与列级血缘是 Mode 集成的核心能力,其提取链路在construct_query_or_dataset与get_upstream_lineage_for_parsed_sql中实现,包含三层 SQL 预处理:
- 定义展开(Definitions):Mode 查询中常使用
{{ @definition as alias }}语法引用公共定义。_replace_definitions会从/definitions端点拉取定义源,递归替换(最大深度 10,检测循环引用),还原成完整 SQL; - 表单模板渲染(Liquid):Mode 查询支持
{% form %}...{% endform %}块与 Liquid 模板变量。normalize_mode_query解析表单参数的默认值并渲染模板,若渲染失败则退回去掉表单块后的 SQL,保证后续解析不因模板语法报错; - 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 → athena、jdbc:bigquery → bigquery、jdbc:postgresql → postgres、jdbc:snowflake → snowflake、jdbc:redshift → redshift等;未识别的前缀使用原始名称并记录 warning; - Snowflake Warehouse 识别:当上游平台为 Snowflake 时,通过正则从 SQL 中提取
use warehouse <name>;语句,将其作为默认数据库名参与血缘解析。
度量/维度标签启发式
启用tag_measures_and_dimensions后,set_field_tags会为 schema 字段打上Dimension或Measure标签。其启发式规则为:字段类型为数值(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),仅供参考