You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

无需Confluent:Kafka Stream对接Solr的Connector与Storm方案咨询

Kafka Stream 与 Solr 集成方案详解

我来帮你拆解这两个问题,给出具体的实现步骤和注意事项,完全避开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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 10:02:59