☰
Apache Pulsar Solr Sink Connector 完全指南:从 Topic 到 Solr 集合的消息落库实践
2026/9/25 1:24:11 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache 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是空字符串根据运行模式有两种写法:
1)SolrCloud 模式:逗号分隔的 Zookeeper 主机列表(可带 chroot),例如localhost:2181,localhost:2182/chroot;
2)Standalone 模式:连接 Solr 的 URL,例如localhost:8983/solr
solrModeString是SolrCloud与 Solr 集群交互时使用的客户端模式,可选值为Standalone与SolrCloud
solrCollectionString是空字符串需要写入记录的 Solr collection 名称
solrCommitWithinMsint否10Solr 更新提交的时间窗口(毫秒)。写入 Solr 的文档会在该时间范围内自动 commit,实现近实时索引
usernameString否空字符串基本认证(Basic Authentication)的用户名。注意:username区分大小写
passwordString否空字符串基本认证的密码。注意:password区分大小写

其中solrUrl、solrMode、solrCollection三个参数在SolrSinkConfig.validate()中被强校验:任一缺失都会抛出NullPointerException(对应消息分别为solrUrl property not set.等);solrCommitWithinMs则要求必须是正整数,否则抛出IllegalArgumentException(solrCommitWithinMs 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": "fake@123" } }

YAML 格式

原文档中的 YAML 示例使用了 JSON 风格的花括号包裹,这里给出标准 YAML 写法(与仓库测试资源 sinkConfig.yaml 保持一致,可直接复制使用):

solrUrl: "localhost:2181,localhost:2182/chroot" solrMode: "SolrCloud" solrCollection: "techproducts" solrCommitWithinMs: 100 username: "fakeuser" password: "fake@123"

需要说明的是,JSON 与 YAML 两种格式的配置内容完全等价:SolrSinkConfig提供了load(String yamlFile)(基于 Jackson YAML 解析)和load(Map<String, Object> map)(基于 Jackson JSON 解析)两个静态加载方法(见 SolrSinkConfig.java),pulsar-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.md):

Flag说明
-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删除连接器。

源码级原理:一条消息如何写入 Solr

1. open:初始化与客户端构建

连接器启动时调用open(Map<String, Object> config, SinkContext sinkContext)(见 SolrAbstractSink.java):

  1. 通过SolrSinkConfig.load(config)解析配置并执行validate()强校验;
  2. 判断username是否为空,决定是否启用 Basic Auth(enableBasicAuth);
  3. 将solrMode转为大写并与枚举比对,非法值直接抛出异常;
  4. 调用getClient(solrMode, solrUrl)构建 Solr 客户端。

客户端构建逻辑(SolrAbstractSink.java)中:

  • Standalone模式使用HttpSolrClient.Builder(url).build(),直接以 HTTP 方式连接单机 Solr;
  • SolrCloud模式使用CloudSolrClient.Builder(zkHosts, chroot).build(),先从 URL 中切分 ZK 主机与 chroot,再基于 Zookeeper 发现集群节点。

2. write:转换、提交与确认

每条消息到达时触发write(Record<T> record)(SolrAbstractSink.java),核心流程为:

  1. 构造UpdateRequest;若solrCommitWithinMs > 0,设置setCommitWithin(...),让 Solr 在该时间窗口内自动提交索引;
  2. 若启用了 Basic Auth,通过setBasicAuthCredentials(username, password)为请求附加认证信息;
  3. 调用抽象方法convert(record)将消息转换为SolrInputDocument并add到请求中;
  4. updateRequest.process(client, solrCollection)将文档提交到目标 collection;
  5. 根据UpdateResponse.getStatus()是否为 0 决定调用record.ack()(确认成功)还是record.fail()(失败重试);捕获到SolrServerException/IOException时同样调用record.fail()并记录告警日志。

3. convert:GenericRecord 到 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.1(solr-solrj),并依赖pulsar-io-core、pulsar-functions-instance与pulsar-client-original,构建时打包为 NAR 归档供 Functions worker 加载。

使用注意事项

  1. 必填项不能缺失:solrUrl、solrMode、solrCollection三者缺失任一都会导致连接器启动失败;
  2. 提交窗口的取舍:solrCommitWithinMs默认 10ms,值越小索引延迟越低,但会带来更频繁的 commit 开销;值越大吞吐更优但搜索可见性滞后,建议按业务对近实时性的要求调整;
  3. 认证凭据区分大小写:username与password均为大小写敏感,配置时需与 Solr 端实际账号完全一致;
  4. 模式与 URL 必须匹配:Standalone模式配 HTTP URL、SolrCloud模式配 ZK 地址列表,交叉配置会导致连接失败;
  5. 消息 schema 与 collection 字段对应:SolrGenericRecordSink按字段名映射,若 Solr collection 中不存在消息中的字段,写入时可能因 schema 校验报错,应提前在 collection 中定义好对应字段。

<输出文章> <输出文章> (重复输出修正)以上为完整正文。 </输出文章>

  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询