无需Confluent:Kafka Stream对接Solr的Connector与Storm方案咨询
我来帮你拆解这两个问题,给出具体的实现步骤和注意事项,完全避开Confluent相关组件:
一、使用 Solr-Kafka Connector(非Confluent版)实现集成
这个连接器基于Apache Kafka Connect开发,无需依赖Confluent平台,完全可以用原生Apache Kafka的Connect组件运行。具体步骤如下:
1. 环境与前置准备
- 部署好Apache Kafka集群(单节点或分布式均可),确保Kafka Connect服务(Apache Kafka自带)可用。
- 部署好Solr实例/集群,提前创建目标数据集合(Collection),并保证集合Schema与Kafka消息字段匹配。
- 安装JDK 8+和Maven(用于编译连接器源码)。
2. 编译连接器Jar包
从该开源连接器的源码仓库获取代码,编译生成可执行Jar包:
# 克隆源码仓库 git clone <该连接器的仓库地址> cd kafka-connect-solr # 编译打包,跳过测试(可选) mvn clean package -DskipTests
编译完成后,将target目录下生成的Jar包放到Kafka Connect的插件目录中(可通过Connect配置的plugin.path参数指定,比如/opt/kafka/plugins)。
3. 配置连接器
创建连接器配置文件(比如solr-sink-connector.properties),填写核心配置项(根据实际环境调整):
# 连接器名称 name=solr-sink-connector # 连接器类路径(需与源码类名一致) connector.class=com.msurendra.kafka.connect.solr.SolrSinkConnector # 并行任务数 tasks.max=1 # 要消费的Kafka主题(多个用逗号分隔,可指定Kafka Stream的输出主题) topics=your-kafka-stream-output-topic # Solr基础URL solr.url=http://localhost:8983/solr # 目标Solr集合名称 solr.collection=your-solr-collection-name # 消息Key/Value转换器(根据数据格式选择,比如JSON) key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter # 是否启用转换器Schema key.converter.schemas.enable=false value.converter.schemas.enable=false # 批量提交到Solr的消息数量 solr.batch.size=100 # Solr请求超时时间(毫秒) solr.request.timeout=5000
注意:如果你的数据来自Kafka Stream处理后的输出,只需将
topics配置为Stream的输出主题即可。
4. 启动Kafka Connect
使用Apache Kafka自带脚本启动Connect(以Standalone模式为例,分布式模式操作类似):
# 进入Kafka安装目录 cd /opt/kafka # 启动Connect,指定通用配置和连接器配置 bin/connect-standalone.sh config/connect-standalone.properties /path/to/solr-sink-connector.properties
启动后查看Connect日志(默认在logs/connect.log),确认连接器无报错、正常运行。
5. 验证集成
往配置的Kafka主题发送测试数据,比如:
bin/kafka-console-producer.sh --broker-list localhost:9092 --topic your-kafka-stream-output-topic # 输入测试JSON数据 {"id":1,"name":"test-stream-data","content":"hello solr from kafka stream"}
登录Solr Admin UI,查看目标集合的文档数量,确认数据成功写入。
二、Apache Storm 接收Kafka Stream数据并清洗后推送至Solr
完全可以实现!Storm天生擅长实时数据流处理,结合Kafka Spout和SolrJ客户端就能完成从消费、清洗到写入Solr的全流程。具体步骤:
1. 依赖准备
在Storm拓扑项目的pom.xml中添加必要依赖:
<!-- Storm Kafka客户端依赖 --> <dependency> <groupId>org.apache.storm</groupId> <artifactId>storm-kafka-client</artifactId> <version>${storm.version}</version> <scope>provided</scope> </dependency> <!-- SolrJ客户端依赖 --> <dependency> <groupId>org.apache.solr</groupId> <artifactId>solr-solrj</artifactId> <version>${solr.version}</version> </dependency>
替换${storm.version}和${solr.version}为你实际使用的版本。
2. 构建Storm拓扑
(1)Kafka Spout配置
创建Kafka Spout消费Kafka Stream的输出主题:
import org.apache.storm.kafka.client.spout.KafkaSpout; import org.apache.storm.kafka.client.spout.KafkaSpoutConfig; // 配置Kafka Spout,直接消费Kafka Stream的输出主题 KafkaSpoutConfig<String, String> spoutConfig = KafkaSpoutConfig.builder("localhost:9092", "your-kafka-stream-output-topic") .setGroupId("storm-kafka-stream-group") .setOffsetCommitPeriodMs(1000) .build(); KafkaSpout<String, String> kafkaSpout = new KafkaSpout<>(spoutConfig);
(2)数据清洗Bolt
自定义Bolt实现数据清洗逻辑(比如过滤无效字段、格式转换):
import org.apache.storm.topology.BasicOutputCollector; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.topology.base.BaseBasicBolt; import org.apache.storm.tuple.Tuple; import org.apache.storm.tuple.Values; public class DataCleaningBolt extends BaseBasicBolt { @Override public void execute(Tuple tuple, BasicOutputCollector collector) { String rawStreamData = tuple.getString(1); // 获取Kafka Stream输出的消息内容 // 自定义清洗逻辑:解析JSON、过滤空值、字段转换等 CleanedData cleanedData = parseAndClean(rawStreamData); // 将清洗后的数据发送到下一个Bolt collector.emit(new Values(cleanedData)); } @Override public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields("cleanedData")); } // 模拟数据清洗方法 private CleanedData parseAndClean(String rawData) { // 实现你的解析与清洗逻辑 return new CleanedData(); } } // 自定义清洗后的数据实体类 class CleanedData { private String id; private String name; private String content; // 省略getter/setter }
(3)Solr写入Bolt
实现Bolt用SolrJ将清洗后的数据写入Solr:
import org.apache.solr.client.solrj.impl.HttpSolrClient; import org.apache.solr.client.solrj.SolrClient; import org.apache.solr.common.SolrInputDocument; import org.apache.storm.topology.BasicOutputCollector; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.topology.base.BaseBasicBolt; import org.apache.storm.tuple.Tuple; public class SolrWriterBolt extends BaseBasicBolt { private transient SolrClient solrClient; @Override public void prepare(java.util.Map conf, org.apache.storm.task.TopologyContext context) { // 初始化Solr客户端 String solrUrl = "http://localhost:8983/solr/your-solr-collection-name"; solrClient = new HttpSolrClient.Builder(solrUrl).build(); } @Override public void execute(Tuple tuple, BasicOutputCollector collector) { CleanedData cleanedData = (CleanedData) tuple.getValueByField("cleanedData"); // 转换为Solr文档 SolrInputDocument doc = new SolrInputDocument(); doc.addField("id", cleanedData.getId()); doc.addField("name", cleanedData.getName()); doc.addField("content", cleanedData.getContent()); // 写入Solr try { solrClient.add(doc); solrClient.commit(); } catch (Exception e) { // 异常处理:重试、日志记录等 e.printStackTrace(); } } @Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // 该Bolt无需输出,无需声明字段 } @Override public void cleanup() { // 关闭Solr客户端 try { solrClient.close(); } catch (Exception e) { e.printStackTrace(); } } }
(4)组装并提交拓扑
将Spout和Bolt组装为完整拓扑并提交:
import org.apache.storm.topology.TopologyBuilder; import org.apache.storm.StormSubmitter; import org.apache.storm.Config; public class KafkaStreamToSolrTopology { public static void main(String[] args) throws Exception { TopologyBuilder builder = new TopologyBuilder(); // 设置Kafka Spout,并行度1 builder.setSpout("kafka-stream-spout", kafkaSpout, 1); // 设置数据清洗Bolt,并行度2 builder.setBolt("data-cleaning-bolt", new DataCleaningBolt(), 2) .shuffleGrouping("kafka-stream-spout"); // 设置Solr写入Bolt,并行度1 builder.setBolt("solr-writer-bolt", new SolrWriterBolt(), 1) .shuffleGrouping("data-cleaning-bolt"); Config config = new Config(); config.setDebug(false); // 提交到Storm集群 StormSubmitter.submitTopology("kafka-stream-to-solr-topology", config, builder.createTopology()); } }
3. 运行与验证
- 打包拓扑Jar包,用Storm命令提交到集群:
storm jar your-topology.jar com.your.package.KafkaStreamToSolrTopology
- 往Kafka Stream的输出主题发送测试数据,查看Solr集合中的数据是否为清洗后的结果,同时通过Storm UI确认拓扑运行正常。
方案对比建议
- Solr-Kafka Connector:适合无需复杂数据处理的场景,配置简单、运维成本低,依托Kafka Connect的容错机制,适合批量或准实时数据同步。
- Apache Storm:适合需要复杂实时计算、多步骤数据清洗、自定义业务逻辑的场景,灵活性更高,但需要编写代码开发拓扑,运维复杂度稍高。
内容的提问来源于stack exchange,提问作者Lijo wilson

