DataHub 集成 Apache Pulsar:基于 Pulsar Admin API 的 Topic 与 Schema 元数据采集指南
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 列表覆盖四类 topicpersistent topics持久化 topicpersistent partitioned topics持久化分区 topicnon-persistent topics非持久化 topicnon-persistent partitioned topics非持久化分区 topic获取每个 topic 的最新 schema/admin/v2/schemas/{tenant}/{namespace}/{topic}/schema数据以tenant / namespace为粒度进行提取topic 连同其可用的 schema 一起被摄取为 DataHub 中的 Dataset 实体。此外schema descriptionschema 描述、schema_versionschema 版本、schema_typeschema 类型以及partitioned是否分区等附加信息会被写入DatasetProperties自定义属性中。概念映射Pulsar 实体到 DataHub 实体的对应关系根据 pulsar 目录下的 READMEPulsar 源在 DataHub 中的实体映射关系如下源概念DataHub 概念说明pulsar平台Data Platform注册为名为pulsar的数据平台Pulsar TopicDataset子类型subType标记为topicPulsar SchemaSchemaField映射到 Avro 或 JSON schema 定义中的字段从源码来看这一映射在 pulsar.py 中以装饰器形式声明platform_name(Pulsar)、config_class(PulsarSourceConfig)并声明了三项能力PLATFORM_INSTANCE默认启用、DOMAINS通过domain配置字段支持、SCHEMA_METADATA默认启用当前支持状态为BETAsupport_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 参数”。快速上手编写摄取 RecipeDataHub 的 Python 摄取工具acryl-datahub通过 RecipeYAML 配置定义 source 与 sink。pulsar源的完整示例位于 pulsar_recipe.ymlsource: 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的完整 URLOAuth 认证方式必填client_idstrNoneOAuth 应用客户端 IDclient_secret密文字符串NoneOAuth 应用客户端密钥token密文字符串None直接使用的访问令牌JWTToken 认证方式必填其中web_service_url有专门的 Pydantic 校验器web_service_url_scheme_host_port见 source_config/pulsar.pyscheme 必须是http或https否则报错Scheme should be http or https, found xxxhostname 会做宽松的 ASCII 校验每段不超过 63 字符、总长不超过 253 字符、不允许非法字符尾部多余的斜杠会被自动清理config_clean.remove_trailing_slashes。单元测试 test_pulsar_source.py 验证了这些行为例如http://localhost:8080/会被规范化为http://localhost:8080而ftp://localhost:8080/与包含非法字符的 hostname 都会被拒绝。认证方式JWT Token 与 OAuthpulsar源支持两种认证方式且二者互斥TokenJWT认证配置token字段即可get_access_token()会直接返回该 JWTOAuth 客户端凭证认证配置issuer_urlclient_idclient_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 模式直接返回配置中的 tokenOAuth 模式则通过 mock 的 OIDC 发现文档与 token 端点换取令牌。采集范围控制参数类型默认值说明tenantsList[str][]要采集的 tenant 列表。留空时需要通过/admin/v2/tenants枚举集群所有 tenant这要求superUser角色显式指定列表时使用 tenant 管理员角色即可tenant_patternsAllowDenyPatternallow.*denypulsartenant 的包含/排除正则默认排除 Pulsar 内部pulsartenantnamespace_patternsAllowDenyPatternallow.*denypublic/functionsnamespace 的包含/排除正则默认排除 Pulsar 内置函数命名空间topic_patternsAllowDenyPatternallow.*deny/__.*$topic 的包含/排除正则默认排除 Pulsar 系统 topicexclude_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上报集群版本请求{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(removedFalse)URL 名称取 topic 完整名如persistent://tenant/namespace/topic并结合platform_instance与env生成 URNschemaMetadata 方面如果 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 完整名得到typepersistent/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、dataAvro/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_versionbroker 版本、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),仅供参考