ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

OpenMetadata NATS 连接器接入指南:JetStream 流摄取、认证与 Schema 管理

2026/9/15 13:03:22 拓冰建站 浏览量
OpenMetadata NATS 连接器接入指南:JetStream 流摄取、认证与 Schema 管理 OpenMetadata NATS 连接器接入指南JetStream 流摄取、认证与 Schema 管理【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata本指南系统讲解 OpenMetadata 中 NATS 消息中间件连接器的完整接入方式涵盖 JetStream 前置要求、连接参数、四种认证方式、TLS 加密配置以及基于 JetStream KV Bucket 的 Schema 摄取机制。读完本文你将能够从零配置一个可用的 NATS 数据源把 JetStream 流以 Topic 的形式接入 OpenMetadata 元数据目录并结合仓库源码理解其底层摄取原理。连接器概述NATS 连接器是 OpenMetadata 消息服务Messaging Service体系中的一员负责将 NATS JetStream 流Stream摄取为 OpenMetadata 中的 Topic 实体。它的定位与其他消息连接器如 Kafka、Kinesis、Redpanda、PubSub一致让消息中间件中的数据结构化、可检索并纳入统一的数据治理体系。从源码结构看NATS 连接器由四部分构成均位于 ingestion/src/metadata/ingestion/source/messaging/nats 目录下service_spec.py连接器注册入口将NatsSource摄取实现与NatsConnection连接处理绑定为 ServiceSpecconnection.py负责建立 NATS 客户端连接、构建 TLS 上下文、执行连接测试metadata.py核心摄取逻辑列出流、抓取消息、解析 Schema 并产出 Topic 实体models.py定义流配置NatsStreamConfig与流状态NatsStreamState的数据模型。连接器的 JSON Schema 定义见 natsConnection.json其中natsServers是唯一必填字段其余字段均为可选。前置要求必须启用 JetStream连接 NATS 服务器前需确保服务器端已启用JetStream这是 NATS 官方提供的持久化消息子系统。启动时可通过两种方式开启命令行方式使用-js标志启动例如nats-server -js配置文件方式在服务器配置文件中设置jetstream: enabled。需要特别注意的是连接器只支持 JetStream 流原因在于 API 层面的硬性限制核心 NATS 主题未启用 JetStream 的普通 subjects无法通过 NATS API 列出因此无法被枚举和摄取。连接器源码中正是通过 JetStream 管理 API$JS.API.STREAM.LIST来获取流列表的见 connection.py这条 API 仅在 JetStream 启用后可用。连接配置详解NATS 连接的配置入口在 OpenMetadata UI 的「添加服务 → Messaging → NATS」页面也可以在 ingestion 工作流 YAML 中以serviceConnection.config方式声明。全部字段及其含义如下。NATS ServersnatsServers必填NATS 服务器地址列表多个地址用逗号分隔。每个 URL 都必须包含协议与端口nats://host1:4222,nats://host2:4222普通连接使用nats://协议TLS 加密连接使用tls://协议例如tls://host1:4222。源码在 connection.py 中按逗号切分该字符串并逐个去除首尾空白后作为servers列表传入nats.connect()若切分后存在空项会直接抛出ValueError因此请勿在列表中写入空元素。Authentication TypeauthType认证方式支持三种任选其一同一时间只能配置一种认证方式认证方式配置字段说明Username and PasswordusernamepasswordNATS 基础认证user:password两字段必须同时提供TokentokenToken 令牌认证仅需一个令牌NKey SeednkeySeedNKey 种子密钥以SU开头基于 NKeys 加密认证当 NATS 服务器允许匿名连接时authType可以留空。JSON Schema 将认证方式建模为authType的oneOf联合类型见 natsConnection.json三种结构互斥且各自声明了必填字段basicAuth必须同时包含username与passwordtokenAuth必须包含tokennkeyAuth必须包含nkeySeed。UsernameusernameNATS 基础认证的用户名。选择「Username and Password」认证方式时与password成对必填。PasswordpasswordNATS 基础认证的密码在 Schema 中标记为format: passwordOpenMetadata 会将其作为敏感信息处理。TokentokenToken 认证使用的令牌字符串选择 Token 认证方式时必填。NKey SeednkeySeedNKey 认证使用的种子密钥字符串以SU开头。NKeys 提供基于 Ed25519 曲线的加密认证能力凭证不会通过网络明文传输——客户端只发送由种子签名的质询响应服务器通过公钥验签完成认证。密钥对可使用官方nkeys工具生成。TLS ConfigurationtlsConfig用于加密 NATS 连接的 TLS/SSL 配置。当服务器启用了 TLS即连接串使用tls://协议时必须配置。TLS 配置在 Schema 中复用verifySSLConfig.json的sslConfig定义包含三个子字段CA CertificatecaCertificatePEM 格式的 CA 证书用于校验 NATS 服务器的 TLS 证书属于最低要求SSL CertificatesslCertificatePEM 格式的客户端证书仅在需要双向 TLSmTLS时才提供SSL KeysslKey与客户端证书配对的 PEM 私钥同样仅在 mTLS 场景下需要。值得注意的是一处源码级约束sslCertificate与sslKey必须成对出现只配置其一会在构建 TLS 上下文时抛出ValueError见 connection.py。从实现细节看证书内容并非直接写入磁盘常驻文件而是由 _write_temp_cert 写入临时.pem文件供ssl.SSLContext加载连接关闭时由_cleanup_temp_certs统一清理避免敏感凭证残留在文件系统中。Additional NATS ConfigadditionalConfig透传给nats.connect()调用的额外客户端配置项为自由格式的键值对象JSON Schema 中additionalProperties: true。可参考 nats.py 客户端文档使用例如设置重连间隔、心跳间隔等高级选项。但以下连接、认证与 TLS 相关选项已被连接器预留不能在additionalConfig中覆盖servers, user, password, token, nkeys_seed, nkeys_seed_str, tls, user_credentials, signature_cb, user_jwt_cb这条约束在 Schema 层natsConnection.json 的propertyNames.not枚举和运行时connection.py 的_RESERVED_CONNECT_OPTIONS集合双重生效即使 Schema 校验被绕过运行期若发现这些保留键也会抛出ValueError。Schema KV BucketschemaKvBucketSchema 摄取功能的开关与数据源。该字段指定一个JetStream KV Bucket的名称Bucket 中按以下约定存放流 SchemaKey必须与流名称Stream Name完全一致Value原始 Schema 文本支持三种格式——Avro根部为{type: record, ...}的 JSONJSON Schema根部包含$schema或properties的 JSONProtobuf以syntax proto3;或syntax proto2;开头的文本。留空则跳过 Schema 摄取。关于格式识别连接器内部实现了 _detect_schema_type先判断是否为 JSON以{开头再依据根部关键字区分 Avrotype为 record/enum/array/fixed与 JSON Schema含$schema或properties非 JSON 文本则根据syntax与message关键字判定 Protobuf否则归为Other。Topic Filter PatterntopicFilterPattern用正则表达式按流名称过滤摄取范围Includes仅摄取匹配这些模式Pattern的流Excludes跳过匹配这些模式的流。两者均留空时摄取全部流。该字段在 Schema 中引用filterPattern.json的filterPattern定义与 OpenMetadata 其他连接器如 Kafka的过滤语义保持一致。完整配置示例以下是仓库中提供的 NATS 摄取工作流完整示例见 ingestion/src/metadata/examples/workflows/nats.yaml包含注释掉的认证、TLS 与 Schema 选项可直接按需取消注释使用source: type: nats serviceName: local_nats serviceConnection: config: type: Nats natsServers: nats://localhost:4222 # Authentication — choose ONE method: # authType: # username: myuser # password: mypassword # Or: # authType: # token: mytoken # Or: # authType: # nkeySeed: SUAA... # TLS (optional): # tlsConfig: # caCertificate: | # -----BEGIN CERTIFICATE----- # sample caCertificateData # -----END CERTIFICATE----- # sslCertificate: | # -----BEGIN CERTIFICATE----- # sample sslCertificateData # -----END CERTIFICATE----- # sslKey: | # -----BEGIN RSA PRIVATE KEY----- # sample sslKeyData # -----END RSA PRIVATE KEY----- # Schema ingestion from JetStream KV bucket (optional): # schemaKvBucket: SCHEMAS sourceConfig: config: type: MessagingMetadata topicFilterPattern: excludes: - _heartbeat.* generateSampleData: true sink: type: metadata-rest config: {} workflowConfig: # loggerLevel: INFO # DEBUG, INFO, WARN or ERROR openMetadataServerConfig: hostPort: http://localhost:8585/api authProvider: openmetadata securityConfig: jwtToken: eyJraWQiOiJHYjM4OWEtOWY3Ni1nZGpzLWE5MmotMDI0MmJrOTQzNTYiLCJ0eXAiOiJKV1QiLCJhbGciOiJSUzI1NiJ9...其中值得注意的实践要点sourceConfig.config.type固定为MessagingMetadata表示这是元数据Metadata类摄取topicFilterPattern.excludes中_heartbeat.*用于跳过系统心跳类流避免噪声数据进入元数据目录generateSampleData: true开启样本数据采集默认开启可通过配置关闭sink.type: metadata-rest将摄取结果写入 OpenMetadata 服务端 REST API。摄取原理从 JetStream 流到 OpenMetadata Topic理解底层实现有助于排障与调优。NATS 连接器继承了消息服务公共基类MessagingServiceSource见 ingestion/src/metadata/ingestion/source/messaging/messaging_service.py整个摄取链路分为三个环节全部实现在 metadata.py 中。1. 流列举分页拉取$JS.API.STREAM.LISTget_topic_list 通过 JetStream 管理 API$JS.API.STREAM.LIST分页枚举流请求体携带{offset: offset}服务端返回total、offset、streams数组连接器严格校验响应offset与请求一致、total完整无missing并保证分页推进next_offset offset任一异常都会抛出NatsApiError每个流的config与state分别解析为NatsStreamConfig与NatsStreamState见 models.py构成 Topic 元数据。2. Topic 实体构建流配置映射为 Topic 属性yield_topic 将流的原生配置映射为 OpenMetadata Topic 实体的标准属性流配置JetStreamTopic 属性OpenMetadatanum_replicas副本数replicationFactormax_msg_size单消息上限maximumMessageSizemax_bytes存储上限retentionSizemax_age保留时长纳秒retentionTime毫秒除以 1,000,000 转换retentionlimits/workqueue/interestcleanupPolicies统一映射为deletesubjects、storage、retention额外写入topicConfig保留原始细节3. 样本数据与 SchemaJetStream 消息读取样本数据当generateSampleData开启时_fetch_sample_messages 从流末尾向前扫描最多 100 条序列_SAMPLE_SCAN_LIMIT抓取最多 10 条_SAMPLE_SIZE、累计不超过 1MB_SAMPLE_BYTE_LIMIT的文本消息作为 Topic 样本消息体为 Base64 编码非 UTF-8 文本如二进制会被跳过。Schema 解析若配置了schemaKvBucket_fetch_schema_from_kv 通过$JS.API.STREAM.MSG.GET.KV_{bucket}按last_by_subj: $KV.{bucket}.{stream_name}拉取对应 Schema 文本Base64 解码后交给schema_parser_config_registry注册的解析器生成schemaFieldsProtobuf 文本会先经merge_and_clean_protobuf_schema预处理再解析。如果流在 Bucket 中无对应消息错误码 10037则跳过该流的 Schema 而不中断整体摄取。连接测试与自动化验证连接器提供两层连接测试便于在创建服务或排障时快速定位问题。测试逻辑实现在 connection.py 的test_connection中包含两个测试步骤GetTopics调用$JS.API.STREAM.LIST验证能否连通服务器并枚举 JetStream 流——未启用 JetStream 的服务器会在此步骤失败CheckSchemaKvBucket仅当配置了schemaKvBucket时执行通过$JS.API.STREAM.INFO.KV_{bucket}验证 KV Bucket 是否可用未配置 Bucket 时抛出SchemaKvBucketNotConfiguredError。这两个步骤可通过 OpenMetadata UI 的「Test」按钮在配置页即时触发也可在自动化工作流Automations Workflow中复用。仓库同时提供了配套测试保障连接器行为单元测试 ingestion/tests/unit/source/messaging/test_nats.py覆盖数据模型解析含额外字段容忍、连接选项构建_build_connect_opts、TLS 上下文构建、Schema 类型探测_detect_schema_type等集成测试目录 ingestion/tests/integration/nats包含populate_nats.py向 NATS 写入测试数据、populate_schemas_kv.py向 KV Bucket 写入 Schema与test_metadata.py端到端摄取验证是理解连接器真实工作方式的最佳参考。注意事项与排障要点JetStream 未启用连接测试「GetTopics」失败、报NatsApiError请确认服务器以-js启动或配置文件包含jetstream: enabled认证互斥三种认证方式只能选择一种UI 上切换认证类型会重置另一类型的表单TLS 证书成对sslCertificate与sslKey必须同时提供否则连接构建报ValueErrorSchema 不生效检查 KV Bucket 中 Key 是否与流名称完全一致大小写敏感且 Value 为支持的三类格式之一未配置schemaKvBucket时 Schema 摄取不会执行保留配置项冲突additionalConfig中不得出现servers、user、token、tls等保留键否则 Schema 校验或运行期检查会拒绝该配置过滤规则topicFilterPattern的 Includes/Excludes 均基于流名称匹配系统类流如_heartbeat.*建议通过 Excludes 排除。小结NATS 连接器为 OpenMetadata 接入 NATS JetStream 生态提供了完整闭环流枚举、Topic 实体映射、样本数据采集与 Schema 解析一应俱全连接层则覆盖了匿名、用户名密码、Token、NKey 四种认证与 TLS/mTLS 加密。配置时只需牢记两个前提——服务器必须启用 JetStreamSchema 摄取依赖约定结构的 KV Bucket——即可快速将 NATS 消息结构纳入 OpenMetadata 的统一治理体系。【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考