尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

Apache Pulsar Solr Sink Connector 实战指南:配置、写入原理与源码解析

发布时间:2026/9/27 21:18:00

资讯中心
01
ARTICLE

Apache Pulsar Solr Sink Connector 实战指南:配置、写入原理与源码解析

Apache Pulsar Solr Sink Connector 实战指南:配置、写入原理与源码解析
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本文以 Apache Pulsar 官方文档 io-solr.md 为核心系统讲解 Solr Sink Connector 的用途、全部配置项、部署方式并结合当前仓库中 pulsar-io/solr 模块的源码与单元测试深入剖析其底层实现原理。读完本文你将掌握如何配置 SolrCloud 与 Standalone 两种模式、各参数的实际作用与默认值、消息是如何被转换为 Solr 文档并落库的以及如何通过测试代码验证整个写入链路。一、Solr Sink Connector 是什么Solr Sink Connector 是 Pulsar IO 框架提供的一个输出端Sink连接器用于从 Pulsar topic 中拉取消息并将消息持久化写入 Solr 的 collection集合。它的典型应用场景是把 Pulsar 中实时产生的业务数据如日志、订单、搜索索引数据持续同步到 Apache Solr供全文检索与聚合分析使用。在 Pulsar IO 体系中Sink 是从 Pulsar 读、向外部系统写的连接器对应的 Source 则是从外部系统读、写入 Pulsar。Solr Sink 属于前者数据流向为Pulsar Topic → Solr Sink Connector → Solr Collection从源码看Solr Sink 的实现位于 pulsar-io/solr 模块Maven artifact 为pulsar-io-solr其类结构如下类职责SolrSinkConfig.java配置模型负责 YAML/Map 加载与参数校验SolrAbstractSink.java抽象基类封装 Solr 客户端创建、写入与 ack/fail 语义SolrGenericRecordSink.java具体实现将GenericRecord转换为SolrInputDocument该模块通过Connector(name solr, type IOType.SINK)注解声明为名为solr的 Sink 连接器见 SolrGenericRecordSink.java依赖 Solr 官方客户端solr-solrj当前仓库中版本为 8.11.1见 pom.xml并使用nifi-nar-maven-plugin打包为 NAR 格式供 Pulsar 加载。二、Sink 配置项详解Solr Sink 的全部配置项及其默认值、是否必填、含义如下表与官方文档一致并补充了源码中的实现细节名称默认值是否必填说明solrUrlnull是SolrCloud 模式下为逗号分隔的 Zookeeper 主机列表可带 chroot如localhost:2181,localhost:2182/chrootStandalone 模式下为 Solr 的连接 URL如localhost:8983/solrsolrModeSolrCloud是与 Solr 集群交互时使用的客户端模式可选值Standalone、SolrCloudsolrCollectionnull是记录要写入的 Solr collection 名称solrCommitWithinMs10否Solr 更新提交的毫秒数未配置时默认 10 msusernamenull否基本认证basic authentication使用的用户名passwordnull否基本认证使用的密码2.1 配置校验规则配置在加载后立即进入校验流程其逻辑定义在 SolrSinkConfig.java 的validate()方法中solrUrl、solrMode、solrCollection三项必须设置否则抛出NullPointerExceptionsolrCommitWithinMs必须为正整数否则抛出IllegalArgumentException消息为 solrCommitWithinMs must be a positive integer.同时solrMode的值在open()时会被转换为大写并与枚举SolrModeSTANDALONE、SOLRCLOUD比对非法值将抛出异常消息为 Illegal Solr mode, valid values are: ...。2.2 配置加载方式从源码看SolrSinkConfig.java 提供了两种加载入口load(String yamlFile)通过 Jackson 的YAMLFactory从 YAML 文件读取配置load(MapString, Object map)从配置 Map 读取Pulsar Sink 运行时正是以 Map 形式下发配置。仓库测试资源 sinkConfig.yaml 给出了一个完整的 JSON 风格 YAML 示例{ solrUrl: localhost:2181,localhost:2182/chroot, solrMode: SolrCloud, solrCollection: techproducts, solrCommitWithinMs: 100, username: fakeuser, password: fake123 }对应的解析与校验测试见 SolrSinkConfigTest.java其中包括正常加载断言solrUrl、solrMode、solrCollection、solrCommitWithinMs、username、password一一比对、缺少solrUrl时的NullPointerException、solrCommitWithinMs为负数时的IllegalArgumentException、非法solrMode值如NotSupport时的IllegalArgumentException以及 chroot 解析测试。三、两种 Solr 模式的连接方式solrMode决定 Sink 使用哪种 Solr 客户端其选择与构建逻辑集中在 SolrAbstractSink.java 的getClient(SolrMode solrMode, String url)静态方法中Standalone 模式单机使用HttpSolrClient.Builder(url)构建客户端solrUrl直接指向 Solr 服务地址例如http://localhost:8983/solr。单元测试 SolrGenericRecordSinkTest.java 中使用的正是solrUrl: http://localhost:8983/solr、solrMode: Standalone的组合。SolrCloud 模式使用CloudSolrClient.Builder(zkHosts, chroot)构建客户端。solrUrl需提供逗号分隔的 Zookeeper 地址列表并可按indexOf(/)切分出 chroot 路径首个/之前的字符串解析为 ZK 主机列表逗号分割其后的路径作为 chroot例如localhost:2181,localhost:2182/chroot会解析出 ZK 主机localhost:2181、localhost:2182与 chroot/chroot。这一解析逻辑在 SolrSinkConfigTest.java 的validZkChrootTest中有明确验证。注意solrUrl的两种写法取决于solrMode——SolrCloud 下应填写 ZK 集群地址可带 chrootStandalone 下应填写 Solr HTTP 服务地址二者不可混用。四、写入流程从 Pulsar 消息到 Solr 文档Solr Sink 的写入生命周期由Sink接口的三个方法驱动实现在 SolrAbstractSink.java 中4.1 open()初始化客户端open(MapString, Object config, SinkContext sinkContext)负责加载并校验配置、解析模式并构建客户端SolrSinkConfig.load(config)解析配置随后调用validate()校验必填项若username非空则启用基本认证enableBasicAuth !Strings.isNullOrEmpty(getUsername())将solrMode转为大写后匹配枚举非法值抛出异常调用getClient(solrMode, solrUrl)构建SolrClient。4.2 write()逐条写入并处理 ack/failwrite(RecordT record)是核心写入逻辑见 SolrAbstractSink.java流程如下构造UpdateRequest若solrCommitWithinMs 0则设置setCommitWithin(solrCommitWithinMs)——这会让 Solr 在指定毫秒内自动提交避免每条消息都触发一次全量 commit从而提升批量写入吞吐若启用了基本认证则通过setBasicAuthCredentials(username, password)附加认证信息调用convert(record)将 Pulsar 记录转换为SolrInputDocument并加入请求执行updateRequest.process(client, solrCollection)将文档写入指定 collection根据响应状态决定消息语义UpdateResponse.getStatus() 0时record.ack()确认成功否则record.fail()若抛出SolrServerException或IOException同样调用record.fail()并记录告警日志。这一 ack/fail 机制保证了 Pulsar 的至少一次at-least-once投递语义与 Solr 写入结果保持一致。4.3 convert()GenericRecord 到 SolrInputDocument具体的文档转换由子类实现。当前仓库提供的 SolrGenericRecordSink.java 面向带 Schema 的消息GenericRecord它遍历记录的Field列表逐个调用record.getField(field)取值并以字段名调用doc.setField(field.getName(), fieldValue)最终生成SolrInputDocument。这意味着写入 Solr 的文档字段名与 Pulsar 消息 Schema 中的字段名一一对应因此在使用前应确保 Pulsar 消息携带 Schema如 Avro且 Schema 字段与 Solr collection 中的字段定义相匹配。4.4 close()释放客户端close()关闭SolrClient连接释放底层资源。五、测试验证写入链路如何被证明可用仓库为 Solr Sink 提供了完整的单元测试支撑SolrServerUtil.java 使用 Jetty 内嵌方式JettySolrRunner在测试中启动 Standalone Solr 实例默认端口 8983上下文/solr用于真实环境下的读写验证SolrGenericRecordSinkTest.java 展示了完整的使用范式定义 POJOFoo含field1、field2两个字段用AvroSchema.of(Foo.class)编码消息配置solrUrlhttp://localhost:8983/solr、solrModeStandalone、solrCollectiontechproducts、solrCommitWithinMs100然后执行sink.open(configs, null)验证初始化成功SolrSinkConfigTest.java 覆盖了配置加载、必填项校验、非法模式/非法毫秒数等边界场景。这些测试同时印证了本文前述的配置格式与写入行为可作为二次开发或故障排查的参考基线。六、部署与使用建议要在 Pulsar 中使用 Solr Sink典型步骤如下以 Pulsar IO 的标准流程为准准备 Solr确认目标 Solr 集群为 Standalone 或 SolrCloud 模式并提前创建好要写入的 collection如techproducts准备 Pulsar 消息确保写入 Pulsar topic 的消息携带 Schema如 Avro字段名与 Solr collection 字段对应编写 Sink 配置按上文表格配置solrUrl、solrMode、solrCollection等参数通过 Pulsar Admin CLI 创建 Sink使用pulsar-admin sinks create命令指定连接器类型为solr、配置文件路径与目标 topic验证向 topic 生产消息检查 Solr collection 中是否出现对应文档。需要提醒的实践要点必填项不能省略solrUrl、solrMode、solrCollection缺失会导致启动失败异常信息与校验规则见 SolrSinkConfig.java合理设置solrCommitWithinMs该值越小提交越频繁、数据可见性越高但写入吞吐会受影响默认 10 ms 适用于多数实时场景测试中常用 100 ms基本认证按需开启仅在username非空时启用SolrAbstractSink.java若 Solr 开启了认证请务必配置用户名密码SolrCloud 的 chroot 写法solrUrl中/后的部分会被识别为 ZK chroot多个 ZK 主机用逗号分隔具体解析逻辑见getClient()与validZkChrootTest测试。七、总结Solr Sink Connector 是 Pulsar 与 Solr 之间的轻量数据管道通过 SolrSinkConfig 声明配置、SolrAbstractSink 封装客户端与写入语义、SolrGenericRecordSink 完成 Schema 到 Solr 文档的转换。理解solrMode两种模式下的solrUrl语义、solrCommitWithinMs的提交时机以及 ack/fail 与写入响应的对应关系即可在实际项目中正确配置并排障。结合 SolrSinkConfigTest.java 与 SolrGenericRecordSinkTest.java 中的边界用例开发者还可以快速搭建本地验证环境复现并确认 Sink 的完整行为。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar HBase Sink Connector 完全指南配置、原理与实战写入Apache Pulsar HBase Sink Connector 完全指南配置、原理与实战写入 本文聚焦 Apache Pulsar 官方提供的 HBas消息队列后端流处理Apache Pulsar RabbitMQ Connector 实战指南Source 与 Sink 配置、部署与源码原理Apache Pulsar RabbitMQ Connector 实战指南Source 与 Sink 配置、部署与源码原理 本指南围绕 Apache Puls消息队列后端流处理Apache Pulsar ElasticSearch Sink Connector 实战指南配置参数、索引策略与写入原理Apache Pulsar ElasticSearch Sink Connector 实战指南配置参数、索引策略与写入原理 ElasticSearch Sin消息队列后端流处理上一篇猫抓Cat-Catch终极指南三步轻松下载网页视频和流媒体资源下一篇JUnit4测试用例优先级UI暗黑模式支持创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。