新闻详情

新闻详情

首页 / 资讯中心 / 详情

Apache Pulsar Solr Sink Connector 完全指南:从 Topic 到 Solr 集合的消息落库实践

发布时间:2026/9/25 1:24:14来源:尧图网络
Apache Pulsar Solr Sink Connector 完全指南:从 Topic 到 Solr 集合的消息落库实践
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Solr sink connector 是 Apache Pulsar 内置的 IO 连接器之一其职责是从 Pulsar Topic 拉取消息并将消息持久化写入 Solr collection。本文以 io-solr-sink.md 为骨架结合pulsar-io/solr模块的源码与测试完整讲解该连接器的全部配置项、JSON/YAML 配置写法、pulsar-admin部署命令以及底层实现原理帮助读者快速落地消息进 Solr的数据管道。连接器概述与适用场景Solr sink connector 实现的是 Pulsar IO 中标准的 Sink 语义连接器作为 Pulsar 消费者订阅一个或多个输入 Topic将每条消息转换为 Solr 的SolrInputDocument再通过 SolrJ 客户端提交到目标 collection。典型的应用场景包括将业务事件流从 Pulsar 实时索引到 Solr供搜索引擎或分析型应用查询用 Pulsar 作为缓冲层解耦消息生产与 Solr 写入借助 Solr 的 near-real-time 提交能力控制索引延迟配合 Pulsar Functions 的至少一次ATLEAST_ONCE处理保证实现生产–消费–索引的可靠链路。在 Pulsar 源码仓库中该连接器的实现位于 pulsar-io/solr由SolrGenericRecordSink通用记录类型 Sink与SolrAbstractSink抽象基类构成并通过Connector(name solr, type IOType.SINK)注册为名为solr的 sink 类型见 SolrGenericRecordSink.java。核心配置参数详解Solr sink connector 的全部配置集中在SolrSinkConfig类中见 SolrSinkConfig.java包括必填项与选填项如下表所示。参数名类型是否必填默认值说明solrUrlString是空字符串根据运行模式有两种写法1SolrCloud 模式逗号分隔的 Zookeeper 主机列表可带 chroot例如localhost:2181,localhost:2182/chroot2Standalone 模式连接 Solr 的 URL例如localhost:8983/solrsolrModeString是SolrCloud与 Solr 集群交互时使用的客户端模式可选值为Standalone与SolrCloudsolrCollectionString是空字符串需要写入记录的 Solr collection 名称solrCommitWithinMsint否10Solr 更新提交的时间窗口毫秒。写入 Solr 的文档会在该时间范围内自动 commit实现近实时索引usernameString否空字符串基本认证Basic Authentication的用户名。注意username区分大小写passwordString否空字符串基本认证的密码。注意password区分大小写其中solrUrl、solrMode、solrCollection三个参数在SolrSinkConfig.validate()中被强校验任一缺失都会抛出NullPointerException对应消息分别为solrUrl property not set.等solrCommitWithinMs则要求必须是正整数否则抛出IllegalArgumentExceptionsolrCommitWithinMs must be a positive integer.。从源码实现看solrMode在连接器启动时会被转换为大写后映射为内部枚举SolrMode.STANDALONE/SolrMode.SOLRCLOUD见 SolrAbstractSink.java因此配置写成SolrCloud、solrcloud均可但不能写成SolrCloud之外的不合法值否则会抛出IllegalArgumentException并提示合法的取值列表。配置文件示例在使用 Solr sink connector 之前需要先通过以下两种方式之一创建配置文件。下面给出文档中的原始示例并附上可直接运行的规范化版本。JSON 格式{ configs: { solrUrl: localhost:2181,localhost:2182/chroot, solrMode: SolrCloud, solrCollection: techproducts, solrCommitWithinMs: 100, username: fakeuser, password: fake123 } }YAML 格式原文档中的 YAML 示例使用了 JSON 风格的花括号包裹这里给出标准 YAML 写法与仓库测试资源 sinkConfig.yaml 保持一致可直接复制使用solrUrl: localhost:2181,localhost:2182/chroot solrMode: SolrCloud solrCollection: techproducts solrCommitWithinMs: 100 username: fakeuser password: fake123需要说明的是JSON 与 YAML 两种格式的配置内容完全等价SolrSinkConfig提供了load(String yamlFile)基于 Jackson YAML 解析和load(MapString, Object map)基于 Jackson JSON 解析两个静态加载方法见 SolrSinkConfig.javapulsar-admin sinks create的--sink-config与--sink-config-file两条路径分别对应这两种加载方式。按运行模式选择 solrUrl 写法SolrCloud 模式默认solrUrl填写 Zookeeper 地址列表。连接器会按第一个/字符拆分 ZK 主机与 chroot例如localhost:2181,localhost:2182/chroot会被解析为主机列表[localhost:2181, localhost:2182]与 chroot/chroot对应测试见 SolrSinkConfigTest.java。Standalone 模式solrUrl填写 HTTP 形式的 Solr 地址例如http://localhost:8983/solr单机测试即使用此写法见 SolrGenericRecordSinkTest.java。部署与运行使用 pulsar-admin 创建 Sink配置文件就绪后可通过pulsar-admin sinks create将 Solr sink 提交到 Pulsar 集群运行pulsar-admin sinks create \ --tenant public \ --namespace default \ --name solr-sink \ --sink-type solr \ --inputs my-topic \ --sink-config-file solr-sink-config.yaml \ --parallelism 1常用参数说明完整选项见 io-cli.mdFlag说明-t,--sink-typeSink 的 connector 类型内置连接器的类型名由pulsar-io.yaml中的name参数决定Solr 对应solr-i,--inputsSink 的输入 Topic多个 Topic 用逗号分隔--nameSink 名称--tenant/--namespaceSink 所属的租户与命名空间--parallelismSink 实例数并行度--sink-config-file指向 YAML 配置文件的路径--sink-config以 key/value 形式直接传入配置与配置文件二选一--processing-guarantees处理保证投递语义可选ATLEAST_ONCE、ATMOST_ONCE、EFFECTIVELY_ONCE实际语义同时依赖 Sink 自身实现--retain-ordering是否按序消费并写入消息Sink 创建后可通过pulsar-admin sinks status查看运行状态、pulsar-admin sinks update更新配置如修改 collection 或提交窗口、pulsar-admin sinks delete删除连接器。源码级原理一条消息如何写入 Solr1. open初始化与客户端构建连接器启动时调用open(MapString, Object config, SinkContext sinkContext)见 SolrAbstractSink.java通过SolrSinkConfig.load(config)解析配置并执行validate()强校验判断username是否为空决定是否启用 Basic AuthenableBasicAuth将solrMode转为大写并与枚举比对非法值直接抛出异常调用getClient(solrMode, solrUrl)构建 Solr 客户端。客户端构建逻辑SolrAbstractSink.java中Standalone模式使用HttpSolrClient.Builder(url).build()直接以 HTTP 方式连接单机 SolrSolrCloud模式使用CloudSolrClient.Builder(zkHosts, chroot).build()先从 URL 中切分 ZK 主机与 chroot再基于 Zookeeper 发现集群节点。2. write转换、提交与确认每条消息到达时触发write(RecordT record)SolrAbstractSink.java核心流程为构造UpdateRequest若solrCommitWithinMs 0设置setCommitWithin(...)让 Solr 在该时间窗口内自动提交索引若启用了 Basic Auth通过setBasicAuthCredentials(username, password)为请求附加认证信息调用抽象方法convert(record)将消息转换为SolrInputDocument并add到请求中updateRequest.process(client, solrCollection)将文档提交到目标 collection根据UpdateResponse.getStatus()是否为 0 决定调用record.ack()确认成功还是record.fail()失败重试捕获到SolrServerException/IOException时同样调用record.fail()并记录告警日志。3. convertGenericRecord 到 Solr 文档SolrGenericRecordSink面向 Pulsar 的GenericRecord带 schema 的消息实现转换SolrGenericRecordSink.java遍历GenericRecord.getFields()中的所有字段将每个字段名与字段值直接映射为SolrInputDocument的同名字段。这意味着消息的 schema 字段与 Solr collection 的 schema 字段应保持对应关系字段名一致时即可完成索引。4. close资源释放连接器停止时调用close()关闭底层SolrClient释放网络与连接池资源SolrAbstractSink.java。测试与验证仓库内的质量保障仓库在 pulsar-io/solr/src/test 下提供了两类测试可作为理解连接器行为的参考配置解析与校验测试SolrSinkConfigTest.java覆盖 YAML 文件加载、Map 加载、合法配置校验以及三类异常场景——缺少solrUrl抛NullPointerException、solrCommitWithinMs为负数抛IllegalArgumentException、solrMode为NotSupport时因枚举不匹配抛IllegalArgumentException。端到端写入测试SolrGenericRecordSinkTest.java通过 SolrServerUtil.java 在 Jetty 上拉起嵌入式单机 Solr端口 8983用 Avro schema 编码一个Foo对象作为消息验证open与write全链路可正常运行。此外连接器模块的依赖配置见 pulsar-io/solr/pom.xml当前基于 SolrJ 8.11.1solr-solrj并依赖pulsar-io-core、pulsar-functions-instance与pulsar-client-original构建时打包为 NAR 归档供 Functions worker 加载。使用注意事项必填项不能缺失solrUrl、solrMode、solrCollection三者缺失任一都会导致连接器启动失败提交窗口的取舍solrCommitWithinMs默认 10ms值越小索引延迟越低但会带来更频繁的 commit 开销值越大吞吐更优但搜索可见性滞后建议按业务对近实时性的要求调整认证凭据区分大小写username与password均为大小写敏感配置时需与 Solr 端实际账号完全一致模式与 URL 必须匹配Standalone模式配 HTTP URL、SolrCloud模式配 ZK 地址列表交叉配置会导致连接失败消息 schema 与 collection 字段对应SolrGenericRecordSink按字段名映射若 Solr collection 中不存在消息中的字段写入时可能因 schema 校验报错应提前在 collection 中定义好对应字段。输出文章 输出文章 重复输出修正以上为完整正文。 /输出文章赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Solr Sink Connector 配置与源码剖析将 Topic 消息持久化到 Solr CollectionApache Pulsar Solr Sink Connector 配置与源码剖析将 Topic 消息持久化到 Solr Collection Solr si消息队列后端流处理Apache Pulsar Kafka Sink Connector 实战指南将 Pulsar Topic 消息桥接到 KafkaApache Pulsar Kafka Sink Connector 实战指南将 Pulsar Topic 消息桥接到 Kafka Kafka Sink Co消息队列后端流处理Apache Pulsar Redis Sink Connector 完全指南将 Topic 消息实时写入 RedisApache Pulsar Redis Sink Connector 完全指南将 Topic 消息实时写入 Redis 本篇技术指南以 Apache Puls消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

源师兄扩展屏图片显示教程:如何用软串口快速输出自定义图像的完整步骤 2026/9/25 2:48:38

源师兄扩展屏图片显示教程:如何用软串口快速输出自定义图像的完整步骤

源师兄扩展屏图片显示教程:如何用软串口快速输出自定义图像的完整步骤 【免费下载链接】software-serial-module 源师兄扩展项目: 软串口模块 | 由源师兄组织创建 项目地址: https://gitcode.com/yuanshixiong/software-serial-module 💡 一句话说…

阅读更多 →
OpenCV 4.8.0与MinGW编译实战:从CMake配置到Qt集成完全指南 2026/9/25 2:48:38

OpenCV 4.8.0与MinGW编译实战:从CMake配置到Qt集成完全指南

简介:从源代码构建OpenCV是许多Windows开发者避开ABI兼容陷阱的通用思路。MSVC与MinGW采用不同的C运行时和链接库格式,官方预编译包无法直接在GCC工具链下使用。通过CMake生成MinGW Makefiles工程,可以控制模块选择、关闭非必要加速项&#x…

阅读更多 →
astron-agent 控制台前端架构指南:基于 Vite + React 18 的 AI Agent 平台工程实践 2026/9/25 2:48:38

astron-agent 控制台前端架构指南:基于 Vite + React 18 的 AI Agent 平台工程实践

人工智能AI AgentAgent 编排RPA后端前端企业应用 【免费下载链接】astron-agent Enterprise-grade, commercial-friendly agentic workflow platform for building next-generation SuperAgents. 项目地址: https://gitcode.com/gh_mirrors/as/astron-agent 点击查看…

阅读更多 →
Ceph RGW SigV4 签名验证缺陷(CVE-2026-54330):预签名 PUT URL 的 x-amz-* 头注入与权限提升深度剖析 2026/9/25 2:48:38

Ceph RGW SigV4 签名验证缺陷(CVE-2026-54330):预签名 PUT URL 的 x-amz-* 头注入与权限提升深度剖析

存储分布式文件系统对象存储后端高可用 【免费下载链接】ceph Ceph is a distributed object, block, and file storage platform 项目地址: https://gitcode.com/gh_mirrors/ce/ceph 点击查看 免费下载 本篇文章基于 Ceph 官方安全公告 doc/security/CVE-2026-54…

阅读更多 →
从Dify到Continue:AI原生应用API编排与跨系统对接实战解析 2026/9/25 2:48:38

从Dify到Continue:AI原生应用API编排与跨系统对接实战解析

把一整套RAG检索、模型调度、工具调用串成一个可复现的接口,这是这两年我做得最多的“杂活”。涉及的方案从自研状态机,到字节的Coze、阿里的百炼、开源的Dify,走了一圈,慢慢摸清了AI原生应用在API编排这条路上的很多暗坑。前两天…

阅读更多 →
OpenClaw 学习系列之七:Gateway 深度解析——WebSocket 长连接与 Session Key 路由的 config.toml 骨架 2026/9/25 2:48:25

OpenClaw 学习系列之七:Gateway 深度解析——WebSocket 长连接与 Session Key 路由的 config.toml 骨架

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉