求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(); } }
一些注意事项
- 连接池优化:生产环境中,Storm的Bolt和Spark的Partition里可以用MongoDB连接池,而不是每次创建单个连接,提升性能。
- 异常处理:示例里没加异常捕获,实际生产要处理Kafka消费异常、MongoDB查询异常,避免任务崩溃。
- 偏移量管理:Storm的KafkaSpout默认会管理偏移量,Spark Streaming里示例用了手动提交,生产环境建议把偏移量存储到可靠地方(比如ZooKeeper或Kafka自身)。
内容的提问来源于stack exchange,提问作者Pedro Kássio
相关产品推荐
相关产品推荐

