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

求Storm Spout从Kafka取数至MongoDB关联Bolt实现数据增强示例(含Spark Streaming)

没问题,我帮你整理了Storm和Spark Streaming两种方案的完整实现示例,涵盖从Kafka读数据、MongoDB查询增强的核心逻辑,你只需要替换配置参数就能快速适配你的环境!

Apache Storm 实现方案

1. 先搞定依赖(Maven)

首先在你的pom.xml里加入这些依赖,确保版本和你的Storm、Kafka、MongoDB环境匹配:

<dependencies>
    <!-- Storm核心依赖 -->
    <dependency>
        <groupId>org.apache.storm</groupId>
        <artifactId>storm-core</artifactId>
        <version>2.4.0</version>
        <scope>provided</scope>
    </dependency>
    <!-- Storm Kafka客户端 -->
    <dependency>
        <groupId>org.apache.storm</groupId>
        <artifactId>storm-kafka-client</artifactId>
        <version>2.4.0</version>
    </dependency>
    <!-- MongoDB Java驱动 -->
    <dependency>
        <groupId>org.mongodb</groupId>
        <artifactId>mongodb-driver-sync</artifactId>
        <version>4.6.1</version>
    </dependency>
</dependencies>

2. Kafka Spout 配置(从Kafka读personID)

Storm官方提供了KafkaSpout,直接用它来消费Kafka主题的消息,这里我们假设Kafka消息的value就是personID:

import org.apache.storm.kafka.spout.KafkaSpout;
import org.apache.storm.kafka.spout.KafkaSpoutConfig;

public class KafkaPersonIdSpout {
    public static KafkaSpout<String, String> createKafkaSpout() {
        KafkaSpoutConfig<String, String> spoutConfig = KafkaSpoutConfig.builder(
                "kafka-broker1:9092,kafka-broker2:9092", // Kafka集群地址
                "your-person-topic" // 要消费的主题
        )
                .setGroupId("storm-kafka-consumer-group") // 消费者组ID
                .setProp(org.apache.kafka.clients.consumer.ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest")
                .build();
        return new KafkaSpout<>(spoutConfig);
    }
}

3. MongoDB 查询增强 Bolt

这个Bolt负责接收Spout传来的personID,去MongoDB查询对应的地址字段,然后输出增强后的(personID, address)数据。注意:MongoClient要在prepare方法初始化,避免每次execute创建连接浪费资源!

import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.BasicOutputCollector;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseBasicBolt;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.Values;
import com.mongodb.client.MongoClients;
import com.mongodb.client.MongoClient;
import com.mongodb.client.MongoCollection;
import org.bson.Document;
import java.util.Map;

public class MongoEnrichmentBolt extends BaseBasicBolt {
    private MongoClient mongoClient;
    private MongoCollection<Document> personCollection;

    @Override
    public void prepare(Map<String, Object> topoConf, TopologyContext context) {
        // 初始化MongoDB连接(替换成你的MongoDB地址)
        mongoClient = MongoClients.create("mongodb://mongo-host:27017");
        personCollection = mongoClient.getDatabase("your-db-name").getCollection("person-collection");
    }

    @Override
    public void execute(Tuple input, BasicOutputCollector collector) {
        // 从Spout拿到personID(假设消息value就是personID)
        String personId = input.getStringByField("value");
        
        // 查询MongoDB对应的地址
        Document personDoc = personCollection.find(new Document("personID", personId)).first();
        String address = personDoc != null ? personDoc.getString("address") : "Unknown";
        
        // 输出增强后的数据
        collector.emit(new Values(personId, address));
    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(new Fields("personID", "address"));
    }

    @Override
    public void cleanup() {
        // 关闭MongoDB连接
        if (mongoClient != null) {
            mongoClient.close();
        }
    }
}

4. 组装并提交Storm Topology

最后把Spout和Bolt组装成Topology,提交到Storm集群:

import org.apache.storm.topology.TopologyBuilder;
import org.apache.storm.StormSubmitter;
import org.apache.storm.Config;

public class PersonEnrichmentTopology {
    public static void main(String[] args) throws Exception {
        TopologyBuilder builder = new TopologyBuilder();
        
        // 设置Spout,并行度设为2
        builder.setSpout("kafka-person-id-spout", KafkaPersonIdSpout.createKafkaSpout(), 2);
        
        // 设置Bolt,并行度设为4,和Spout做shuffle分组
        builder.setBolt("mongo-enrichment-bolt", new MongoEnrichmentBolt(), 4)
                .shuffleGrouping("kafka-person-id-spout");
        
        Config config = new Config();
        config.setDebug(false);
        
        // 提交到Storm集群(本地测试可以用LocalCluster)
        StormSubmitter.submitTopology("person-enrichment-topology", config, builder.createTopology());
    }
}

Spark Streaming 实现方案

如果用Spark Streaming来做,代码会更简洁,而且可以利用Spark的分布式能力,这里用Spark Streaming + Kafka 0.10集成,配合MongoDB Spark Connector来实现:

1. 依赖配置(Maven)

<dependencies>
    <!-- Spark Streaming核心 -->
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-streaming_2.12</artifactId>
        <version>3.3.0</version>
    </dependency>
    <!-- Spark Streaming Kafka集成 -->
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-streaming-kafka-0-10_2.12</artifactId>
        <version>3.3.0</version>
    </dependency>
    <!-- MongoDB Spark Connector -->
    <dependency>
        <groupId>org.mongodb.spark</groupId>
        <artifactId>mongodb-spark-connector_2.12</artifactId>
        <version>10.1.1</version>
    </dependency>
</dependencies>

2. 核心代码实现

这里我们用foreachRDD来处理每个批次的Kafka数据,然后通过MongoDB Connector查询地址,注意用foreachPartition减少MongoDB连接数:

import org.apache.spark.SparkConf;
import org.apache.spark.streaming.Durations;
import org.apache.spark.streaming.api.java.JavaDStream;
import org.apache.spark.streaming.api.java.JavaStreamingContext;
import org.apache.spark.streaming.kafka010.ConsumerStrategies;
import org.apache.spark.streaming.kafka010.KafkaUtils;
import org.apache.spark.streaming.kafka010.LocationStrategies;
import com.mongodb.client.MongoClients;
import com.mongodb.client.MongoClient;
import com.mongodb.client.MongoCollection;
import org.bson.Document;
import java.util.*;

public class SparkPersonEnrichment {
    public static void main(String[] args) throws InterruptedException {
        // 初始化Spark配置
        SparkConf conf = new SparkConf()
                .setAppName("PersonDataEnrichment")
                .setMaster("local[*]"); // 本地测试用,集群部署去掉这个
        
        // 初始化StreamingContext,批次间隔5秒
        JavaStreamingContext jssc = new JavaStreamingContext(conf, Durations.seconds(5));
        
        // Kafka配置参数
        Map<String, Object> kafkaParams = new HashMap<>();
        kafkaParams.put("bootstrap.servers", "kafka-broker1:9092,kafka-broker2:9092");
        kafkaParams.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        kafkaParams.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        kafkaParams.put("group.id", "spark-kafka-consumer-group");
        kafkaParams.put("auto.offset.reset", "latest");
        kafkaParams.put("enable.auto.commit", false);
        
        // 要消费的Kafka主题
        Collection<String> topics = Arrays.asList("your-person-topic");
        
        // 创建Kafka DStream
        JavaDStream<String> personIdStream = KafkaUtils.createDirectStream(
                jssc,
                LocationStrategies.PreferConsistent(),
                ConsumerStrategies.Subscribe(topics, kafkaParams)
        ).map(record -> record.value()); // 提取消息的value(即personID)
        
        // 数据增强:查询MongoDB获取地址
        personIdStream.foreachRDD(rdd -> {
            rdd.foreachPartition(partition -> {
                // 每个Partition初始化一次MongoDB连接,避免频繁创建
                MongoClient mongoClient = MongoClients.create("mongodb://mongo-host:27017");
                MongoCollection<Document> personCollection = mongoClient.getDatabase("your-db-name").getCollection("person-collection");
                
                partition.forEachRemaining(personId -> {
                    Document personDoc = personCollection.find(new Document("personID", personId)).first();
                    String address = personDoc != null ? personDoc.getString("address") : "Unknown";
                    // 这里可以把增强后的数据输出到其他存储,或者打印测试
                    System.out.println("Enhanced Data: personID=" + personId + ", address=" + address);
                });
                
                mongoClient.close();
            });
            // 手动提交Kafka偏移量
            ((org.apache.spark.streaming.kafka010.HasOffsetRanges) rdd.rdd()).offsetRanges();
        });
        
        // 启动StreamingContext
        jssc.start();
        jssc.awaitTermination();
    }
}

一些注意事项

  1. 连接池优化:生产环境中,Storm的Bolt和Spark的Partition里可以用MongoDB连接池,而不是每次创建单个连接,提升性能。
  2. 异常处理:示例里没加异常捕获,实际生产要处理Kafka消费异常、MongoDB查询异常,避免任务崩溃。
  3. 偏移量管理:Storm的KafkaSpout默认会管理偏移量,Spark Streaming里示例用了手动提交,生产环境建议把偏移量存储到可靠地方(比如ZooKeeper或Kafka自身)。

内容的提问来源于stack exchange,提问作者Pedro Kássio

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:45:14