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 实例中提取topic与schema元数据,并写入 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 Topic | Dataset | 子类型(subType)标记为topic |
| Pulsar Schema | SchemaField | 映射到 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_url | str | http://localhost:8080 | Pulsar 集群的 Web 服务地址(即 Admin REST API 地址) |
timeout | int | 5 | 等待 Pulsar REST API 返回数据的超时时间(秒) |
verify_ssl | bool / str | true | 布尔值表示是否校验服务器 TLS 证书;字符串表示 CA bundle 的文件路径 |
issuer_url | str | None | 自定义授权服务器(Authorization Server)的完整 URL,OAuth 认证方式必填 |
client_id | str | None | OAuth 应用客户端 ID |
client_secret | 密文字符串 | None | OAuth 应用客户端密钥 |
token | 密文字符串 | None | 直接使用的访问令牌(JWT),Token 认证方式必填 |
其中web_service_url有专门的 Pydantic 校验器(web_service_url_scheme_host_port,见 source_config/pulsar.py):
- scheme 必须是
http或https,否则报错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源支持两种认证方式,且二者互斥:
- Token(JWT)认证:配置
token字段即可,get_access_token()会直接返回该 JWT; - 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)确保:
token与issuer_url不能同时设置(报错Expected only one authentication method, either issuer_url or token.);- 设置了
issuer_url则client_id与client_secret必须同时提供。
这两条路径在 test_pulsar_source.py 中均有测试覆盖:JWT 模式直接返回配置中的 token,OAuth 模式则通过 mock 的 OIDC 发现文档与 token 端点换取令牌。
采集范围控制
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
tenants | List[str] | [] | 要采集的 tenant 列表。留空时需要通过/admin/v2/tenants枚举集群所有 tenant,这要求superUser角色;显式指定列表时使用 tenant 管理员角色即可 |
tenant_patterns | AllowDenyPattern | allow.*,denypulsar | tenant 的包含/排除正则,默认排除 Pulsar 内部pulsartenant |
namespace_patterns | AllowDenyPattern | allow.*,denypublic/functions | namespace 的包含/排除正则,默认排除 Pulsar 内置函数命名空间 |
topic_patterns | AllowDenyPattern | allow.*,deny/__.*$ | topic 的包含/排除正则,默认排除 Pulsar 系统 topic |
exclude_individual_partitions | bool | true | 是否排除分区 topic 的各个独立 partition。关闭时一个 100 分区的 topic 会生成 100 个 Dataset 实体 |
domain | Dict[str, AllowDenyPattern] | {} | 按 topic 完整名称匹配并附加 DataHub domain 的规则 |
stateful_ingestion | 配置对象 | None | 有状态摄取配置(过期实体清理) |
oid_config | dict | {} | 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):
上报集群版本:请求
{base_url}/brokers/version并将 Pulsar broker 版本写入报告对象PulsarSourceReport;枚举 tenants:若未显式配置
tenants,则请求{base_url}/tenants(需要superUser);随后逐 tenant 用tenant_patterns过滤;枚举 namespaces:对每个通过过滤的 tenant,请求
{base_url}/namespaces/{tenant}(tenant 管理员角色即可完成),再用namespace_patterns过滤;枚举 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属性做准备;逐 topic 提取:通过
topic_patterns过滤后调用_extract_record()生成元数据工作单元(MetadataWorkUnit)。
在_extract_record()中(pulsar.py),每个 topic 依次产出如下工作单元:
- Dataset 实体:
StatusClass(removed=False),URL 名称取 topic 完整名(如persistent://tenant/namespace/topic),并结合platform_instance与env生成 URN; - schemaMetadata 方面:如果 topic 存在 schema,则生成
SchemaMetadata,其中schemaName为 Avro 的namespace.name全限定名,version取 Pulsar schema 版本,hash为 schema 字符串的 MD5,平台 schema 以KafkaSchema结构承载原始 schema 文本与类型; - datasetProperties 方面:将 Pulsar schema 自带 properties 合并
schema_version、schema_type、partitioned三个字段写入customProperties,同时把 schema 的doc字段作为 Dataset 描述; - browsePaths 方面:路径格式为
/{env}/{platform}/{platform_instance}/{tenant}/{namespace}/{topic}(未配置platform_instance时省略该段),便于在 DataHub 中按路径浏览; - dataPlatformInstance 方面:配置了
platform_instance时输出,用于多实例区分; - subTypes 方面:标记
TOPIC子类型; - domains 方面:若配置了
domain且 topic 完整名匹配对应 pattern,则附加相应 domain。
Topic 与 Schema 的解析细节
- Topic 解析:
PulsarTopic类按[: /]分割 topic 完整名,得到type(persistent/non-persistent)、tenant、namespace、topic四个组成部分(pulsar.py)。单元测试test_pulsar_source_parse_topic_string验证了persistent://tenant/namespace/topic的解析结果(test_pulsar_source.py)。 - Schema 解析:
PulsarSchema类从 schema 响应中提取version、data(Avro/JSON 文本)、type、properties,并解析出全限定 schema 名(namespace + "." + name)与描述(doc字段)。schema 数据为空时会回退到空对象并记录告警,JSON 解析失败也会被捕获并记日志(pulsar.py)。 - 字段级解析:只有当 schema 类型为
AVRO或JSON时,才调用schema_util.avro_schema_to_mce_fields()把 schema 文本转换为SchemaField列表;其他类型(如PROTOBUF、STRING等)当前未实现字段解析,会记录一条警告并跳过(pulsar.py)。 - 无 schema 的 topic:请求 schema 时若返回 404,说明该 topic 要么没有 schema、要么尚无消息写入,插件会记录
NoSchemaFound类警告并继续处理下一个 topic,不会中断整体摄取(pulsar.py)。
摄取报告
运行状态由 source_report/pulsar.py 中的PulsarSourceReport统计,包括:pulsar_version(broker 版本)、tenants_scanned、namespaces_scanned、topics_scanned计数,以及被过滤掉的 tenant / namespace / topic 列表(tenants_filtered、namespaces_filtered、topics_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)。
故障排查
如果摄取失败,建议按以下顺序排查:
- 校验凭证与权限:确认 token 或 OAuth 配置有效、账号角色满足要求(枚举全部 tenant 需要
superUser;指定tenants列表时可用 tenant 管理员角色); - 校验连通性:确认
web_service_url可达、端口开放、verify_ssl与 TLS 证书匹配; - 校验范围过滤:确认
tenant_patterns、namespace_patterns、topic_patterns与exclude_individual_partitions没有误伤目标实体; - 检查摄取日志:重点关注
PulsarSourceReport中的告警与统计(如NoSchemaFound、HTTPError、tenants_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),仅供参考