DataHub 集成 Apache Pulsar:基于 Pulsar Admin API 的 Topic 与 Schema 元数据采集指南
2026/9/20 11:44:20 网站建设 项目流程

DataHub 集成 Apache Pulsar:基于 Pulsar Admin API 的 Topic 与 Schema 元数据采集指南

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

导读

本文围绕 DataHub 元数据摄取框架中的pulsar模块展开,讲解如何将 Apache Pulsar 实例中的 topic 与 schema 元数据采集并注册为 DataHub 的 Dataset 实体。读完本文,你将掌握pulsar源的配置方法、JWT / OAuth 两种认证接入方式、tenant–namespace–topic 三级扫描流程、schema 字段解析原理,以及生产环境下的安全建议与常见故障排查手段。

模块概览:pulsar源能做什么

pulsar模块面向生产级摄取(production ingestion)工作流,其核心职责是从 Apache Pulsar 实例中提取topicschema元数据,并写入 DataHub(参见 pulsar_pre.md)。该插件通过调用 Pulsar 官方提供的Admin REST API与 Pulsar 实例交互,具体使用到以下端点:

  • 获取已有 tenant 列表(/admin/v2/tenants
  • 获取每个 tenant 关联的 namespace 列表(/admin/v2/namespaces/{tenant}
  • 获取每个 namespace 关联的 topic 列表,覆盖四类 topic:
    • persistent topics(持久化 topic)
    • persistent partitioned topics(持久化分区 topic)
    • non-persistent topics(非持久化 topic)
    • non-persistent partitioned topics(非持久化分区 topic)
  • 获取每个 topic 的最新 schema(/admin/v2/schemas/{tenant}/{namespace}/{topic}/schema

数据以tenant / namespace为粒度进行提取,topic 连同其可用的 schema 一起被摄取为 DataHub 中的 Dataset 实体。此外,schema description(schema 描述)、schema_version(schema 版本)、schema_type(schema 类型)以及partitioned(是否分区)等附加信息会被写入DatasetProperties自定义属性中。

概念映射:Pulsar 实体到 DataHub 实体的对应关系

根据 pulsar 目录下的 README,Pulsar 源在 DataHub 中的实体映射关系如下:

源概念DataHub 概念说明
pulsar(平台)Data Platform注册为名为pulsar的数据平台
Pulsar TopicDataset子类型(subType)标记为topic
Pulsar SchemaSchemaField映射到 Avro 或 JSON schema 定义中的字段

从源码来看,这一映射在 pulsar.py 中以装饰器形式声明:@platform_name("Pulsar")@config_class(PulsarSourceConfig),并声明了三项能力:PLATFORM_INSTANCE(默认启用)、DOMAINS(通过domain配置字段支持)、SCHEMA_METADATA(默认启用),当前支持状态为BETA@support_status(SupportStatus.BETA))。摄取出的每个 Dataset 实体都会打上SubTypesClass(typeNames=[DatasetSubTypes.TOPIC]),便于在 DataHub UI 中识别为“Topic”类型实体(见 pulsar.py)。

前置条件

根据 pulsar_pre.md,接入前需要满足:

  • Pulsar 实例:可访问的 Pulsar 集群,若已开启认证则需要有效的访问令牌(access token);
  • 版本要求:Pulsar 2.7.0 或更高版本;
  • 角色权限:需要superUser角色才能列出实例内的所有 tenant。

提示:列出 Pulsar 实例内全部现有 tenant 需要superUser角色。如果只有 tenant 管理员(tenant admin)权限,可以在配置中显式指定要采集的 tenant 列表,这样无需superUser也能正常工作(详见下文“tenants 参数”)。

快速上手:编写摄取 Recipe

DataHub 的 Python 摄取工具(acryl-datahub)通过 Recipe(YAML 配置)定义 source 与 sink。pulsar源的完整示例位于 pulsar_recipe.yml:

source: type: "pulsar" config: env: "TEST" platform_instance: "local" ## Pulsar client connection config ## web_service_url: "https://localhost:8443" verify_ssl: "/opt/certs/ca.cert.pem" # Issuer url for auth document, for example "http://localhost:8083/realms/pulsar" issuer_url: <issuer_url> client_id: ${CLIENT_ID} client_secret: ${CLIENT_SECRET} # Tenant list to scrape tenants: - tenant_1 - tenant_2 # Topic filter pattern topic_patterns: allow: - ".*sales.*" sink: # sink configs

运行摄取命令:

datahub ingest -c pulsar_recipe.yml

示例中的${CLIENT_ID}${CLIENT_SECRET}使用环境变量替换,避免把敏感信息直接写入配置文件。

配置参数详解

所有配置项都在 source_config/pulsar.py 中通过 Pydantic 模型PulsarSourceConfig定义。该模型同时继承自StatefulIngestionConfigBase(支持有状态摄取/过期实体清理)、PlatformInstanceConfigMixin(支持多平台实例)与EnvConfigMixin(支持环境标注)。

连接与安全相关

参数类型默认值说明
web_service_urlstrhttp://localhost:8080Pulsar 集群的 Web 服务地址(即 Admin REST API 地址)
timeoutint5等待 Pulsar REST API 返回数据的超时时间(秒)
verify_sslbool / strtrue布尔值表示是否校验服务器 TLS 证书;字符串表示 CA bundle 的文件路径
issuer_urlstrNone自定义授权服务器(Authorization Server)的完整 URL,OAuth 认证方式必填
client_idstrNoneOAuth 应用客户端 ID
client_secret密文字符串NoneOAuth 应用客户端密钥
token密文字符串None直接使用的访问令牌(JWT),Token 认证方式必填

其中web_service_url有专门的 Pydantic 校验器(web_service_url_scheme_host_port,见 source_config/pulsar.py):

  • scheme 必须是httphttps,否则报错Scheme should be http or https, found xxx
  • hostname 会做宽松的 ASCII 校验(每段不超过 63 字符、总长不超过 253 字符、不允许非法字符);
  • 尾部多余的斜杠会被自动清理(config_clean.remove_trailing_slashes)。

单元测试 test_pulsar_source.py 验证了这些行为:例如http://localhost:8080/会被规范化为http://localhost:8080,而ftp://localhost:8080/与包含非法字符&的 hostname 都会被拒绝。

认证方式:JWT Token 与 OAuth

pulsar源支持两种认证方式,且二者互斥

  1. Token(JWT)认证:配置token字段即可,get_access_token()会直接返回该 JWT;
  2. OAuth 客户端凭证认证:配置issuer_url+client_id+client_secret,插件会先请求{issuer_url}/.well-known/openid-configuration发现文档以获取token_endpoint,再以client_credentials模式换取access_token(详见 pulsar.py)。

配置模型中的两个校验器(见 source_config/pulsar.py)确保:

  • tokenissuer_url不能同时设置(报错Expected only one authentication method, either issuer_url or token.);
  • 设置了issuer_urlclient_idclient_secret必须同时提供。

这两条路径在 test_pulsar_source.py 中均有测试覆盖:JWT 模式直接返回配置中的 token,OAuth 模式则通过 mock 的 OIDC 发现文档与 token 端点换取令牌。

采集范围控制

参数类型默认值说明
tenantsList[str][]要采集的 tenant 列表。留空时需要通过/admin/v2/tenants枚举集群所有 tenant,这要求superUser角色;显式指定列表时使用 tenant 管理员角色即可
tenant_patternsAllowDenyPatternallow.*,denypulsartenant 的包含/排除正则,默认排除 Pulsar 内部pulsartenant
namespace_patternsAllowDenyPatternallow.*,denypublic/functionsnamespace 的包含/排除正则,默认排除 Pulsar 内置函数命名空间
topic_patternsAllowDenyPatternallow.*,deny/__.*$topic 的包含/排除正则,默认排除 Pulsar 系统 topic
exclude_individual_partitionsbooltrue是否排除分区 topic 的各个独立 partition。关闭时一个 100 分区的 topic 会生成 100 个 Dataset 实体
domainDict[str, AllowDenyPattern]{}按 topic 完整名称匹配并附加 DataHub domain 的规则
stateful_ingestion配置对象None有状态摄取配置(过期实体清理)
oid_configdict{}OpenID 发现文档的占位容器,一般无需手动配置

其中exclude_individual_partitions的实现值得注意:在create()工厂方法中,若该开关为true,插件会自动向topic_patterns.deny追加一条r".*-partition-[0-9]+"正则,从而把形如xxx-partition-0的单个分区过滤掉(见 pulsar.py)。

底层原理:三层扫描与四端点 topic 枚举

pulsar源的核心摄取流程由get_workunits_internal()驱动(见 pulsar.py):

  1. 上报集群版本:请求{base_url}/brokers/version并将 Pulsar broker 版本写入报告对象PulsarSourceReport

  2. 枚举 tenants:若未显式配置tenants,则请求{base_url}/tenants(需要superUser);随后逐 tenant 用tenant_patterns过滤;

  3. 枚举 namespaces:对每个通过过滤的 tenant,请求{base_url}/namespaces/{tenant}(tenant 管理员角色即可完成),再用namespace_patterns过滤;

  4. 枚举 topics:对每个 namespace 依次请求以下四个端点(对应上文四类 topic):

    • {base_url}/persistent/{tenant}/{namespace}
    • {base_url}/persistent/{tenant}/{namespace}/partitioned
    • {base_url}/non-persistent/{tenant}/{namespace}
    • {base_url}/non-persistent/{tenant}/{namespace}/partitioned

    端点的 URL 以/partitioned结尾即表示该批 topic 是分区 topic,插件据此构造{topic: partitioned}映射,为后续写入partitioned属性做准备;

  5. 逐 topic 提取:通过topic_patterns过滤后调用_extract_record()生成元数据工作单元(MetadataWorkUnit)。

_extract_record()中(pulsar.py),每个 topic 依次产出如下工作单元:

  1. Dataset 实体StatusClass(removed=False),URL 名称取 topic 完整名(如persistent://tenant/namespace/topic),并结合platform_instanceenv生成 URN;
  2. schemaMetadata 方面:如果 topic 存在 schema,则生成SchemaMetadata,其中schemaName为 Avro 的namespace.name全限定名,version取 Pulsar schema 版本,hash为 schema 字符串的 MD5,平台 schema 以KafkaSchema结构承载原始 schema 文本与类型;
  3. datasetProperties 方面:将 Pulsar schema 自带 properties 合并schema_versionschema_typepartitioned三个字段写入customProperties,同时把 schema 的doc字段作为 Dataset 描述;
  4. browsePaths 方面:路径格式为/{env}/{platform}/{platform_instance}/{tenant}/{namespace}/{topic}(未配置platform_instance时省略该段),便于在 DataHub 中按路径浏览;
  5. dataPlatformInstance 方面:配置了platform_instance时输出,用于多实例区分;
  6. subTypes 方面:标记TOPIC子类型;
  7. domains 方面:若配置了domain且 topic 完整名匹配对应 pattern,则附加相应 domain。

Topic 与 Schema 的解析细节

  • Topic 解析PulsarTopic类按[: /]分割 topic 完整名,得到type(persistent/non-persistent)、tenantnamespacetopic四个组成部分(pulsar.py)。单元测试test_pulsar_source_parse_topic_string验证了persistent://tenant/namespace/topic的解析结果(test_pulsar_source.py)。
  • Schema 解析PulsarSchema类从 schema 响应中提取versiondata(Avro/JSON 文本)、typeproperties,并解析出全限定 schema 名(namespace + "." + name)与描述(doc字段)。schema 数据为空时会回退到空对象并记录告警,JSON 解析失败也会被捕获并记日志(pulsar.py)。
  • 字段级解析:只有当 schema 类型为AVROJSON时,才调用schema_util.avro_schema_to_mce_fields()把 schema 文本转换为SchemaField列表;其他类型(如PROTOBUFSTRING等)当前未实现字段解析,会记录一条警告并跳过(pulsar.py)。
  • 无 schema 的 topic:请求 schema 时若返回 404,说明该 topic 要么没有 schema、要么尚无消息写入,插件会记录NoSchemaFound类警告并继续处理下一个 topic,不会中断整体摄取(pulsar.py)。

摄取报告

运行状态由 source_report/pulsar.py 中的PulsarSourceReport统计,包括:pulsar_version(broker 版本)、tenants_scannednamespaces_scannedtopics_scanned计数,以及被过滤掉的 tenant / namespace / topic 列表(tenants_filterednamespaces_filteredtopics_filtered)。这些信息会随摄取结束输出到日志,可用于核对采集范围是否符合预期。

能力、限制与故障排查

pulsar_post.md 对模块行为给出了明确的边界说明:

能力

  • 支持平台实例(platform instance)、domain 附加与 schema 元数据摄取(见上文装饰器声明的能力列表);
  • 支持有状态摄取与过期实体清理(stateful_ingestion配置项,README 中也提到该集成涵盖 topic、connector、pipeline、job 等流式实体并支持状态化删除检测)。

生产环境安全建议

生产环境必须始终启用 TLS 加密,并使用变量替换(variable substitution)处理敏感信息(例如${CLIENT_ID}${CLIENT_SECRET})。

对应到配置中即:web_service_url使用https://协议,verify_ssl指向受信任的 CA bundle 路径,token/client_secret等敏感字段通过环境变量注入。

限制

模块行为受 Pulsar 源 API、权限以及平台暴露的元数据范围约束,例如:非AVRO/JSON类型的 schema 暂不支持字段级解析;topic 自带的自定义属性(Pulsar 2.10.0 特性)在源码中仍以# TODO Add topic properties (Pulsar 2.10.0 feature)标注,尚未写入 DatasetProperties(见 pulsar.py)。

故障排查

如果摄取失败,建议按以下顺序排查:

  1. 校验凭证与权限:确认 token 或 OAuth 配置有效、账号角色满足要求(枚举全部 tenant 需要superUser;指定tenants列表时可用 tenant 管理员角色);
  2. 校验连通性:确认web_service_url可达、端口开放、verify_ssl与 TLS 证书匹配;
  3. 校验范围过滤:确认tenant_patternsnamespace_patternstopic_patternsexclude_individual_partitions没有误伤目标实体;
  4. 检查摄取日志:重点关注PulsarSourceReport中的告警与统计(如NoSchemaFoundHTTPErrortenants_scanned等),根据 source 特有错误信息调整配置。

相关文档

  • pulsar_pre.md(概览与前置条件)
  • pulsar_post.md(能力、限制与故障排查)
  • pulsar_recipe.yml(完整摄取示例)
  • pulsar 源 README(概念映射)
  • pulsar.py(源实现)
  • PulsarSourceConfig(配置模型)
  • PulsarSourceReport(摄取报告)
  • test_pulsar_source.py(单元测试)

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

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

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

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

立即咨询