- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
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),包括必填项与选填项,如下表所示。
| 参数名 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
solrUrl | String | 是 | 空字符串 | 根据运行模式有两种写法: 1)SolrCloud 模式:逗号分隔的 Zookeeper 主机列表(可带 chroot),例如 localhost:2181,localhost:2182/chroot;2)Standalone 模式:连接 Solr 的 URL,例如 localhost:8983/solr |
solrMode | String | 是 | SolrCloud | 与 Solr 集群交互时使用的客户端模式,可选值为Standalone与SolrCloud |
solrCollection | String | 是 | 空字符串 | 需要写入记录的 Solr collection 名称 |
solrCommitWithinMs | int | 否 | 10 | Solr 更新提交的时间窗口(毫秒)。写入 Solr 的文档会在该时间范围内自动 commit,实现近实时索引 |
username | String | 否 | 空字符串 | 基本认证(Basic Authentication)的用户名。注意:username区分大小写 |
password | String | 否 | 空字符串 | 基本认证的密码。注意: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-type | Sink 的 connector 类型,内置连接器的类型名由pulsar-io.yaml中的name参数决定,Solr 对应solr |
-i,--inputs | Sink 的输入 Topic(多个 Topic 用逗号分隔) |
--name | Sink 名称 |
--tenant/--namespace | Sink 所属的租户与命名空间 |
--parallelism | Sink 实例数(并行度) |
--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):
- 通过
SolrSinkConfig.load(config)解析配置并执行validate()强校验; - 判断
username是否为空,决定是否启用 Basic Auth(enableBasicAuth); - 将
solrMode转为大写并与枚举比对,非法值直接抛出异常; - 调用
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),核心流程为:
- 构造
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. 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 加载。
使用注意事项
- 必填项不能缺失:
solrUrl、solrMode、solrCollection三者缺失任一都会导致连接器启动失败; - 提交窗口的取舍:
solrCommitWithinMs默认 10ms,值越小索引延迟越低,但会带来更频繁的 commit 开销;值越大吞吐更优但搜索可见性滞后,建议按业务对近实时性的要求调整; - 认证凭据区分大小写:
username与password均为大小写敏感,配置时需与 Solr 端实际账号完全一致; - 模式与 URL 必须匹配:
Standalone模式配 HTTP URL、SolrCloud模式配 ZK 地址列表,交叉配置会导致连接失败; - 消息 schema 与 collection 字段对应:
SolrGenericRecordSink按字段名映射,若 Solr collection 中不存在消息中的字段,写入时可能因 schema 校验报错,应提前在 collection 中定义好对应字段。
<输出文章> <输出文章> (重复输出修正)以上为完整正文。 </输出文章>
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar Solr Sink Connector 配置与源码剖析:将 Topic 消息持久化到 Solr Collection
Apache Pulsar Solr Sink Connector 配置与源码剖析:将 Topic 消息持久化到 Solr Collection Solr si
消息队列后端流处理Apache Pulsar Kafka Sink Connector 实战指南:将 Pulsar Topic 消息桥接到 Kafka
Apache Pulsar Kafka Sink Connector 实战指南:将 Pulsar Topic 消息桥接到 Kafka Kafka Sink Co
消息队列后端流处理Apache Pulsar Redis Sink Connector 完全指南:将 Topic 消息实时写入 Redis
Apache Pulsar Redis Sink Connector 完全指南:将 Topic 消息实时写入 Redis 本篇技术指南以 Apache Puls
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考