Java 8+Storm 1.1实现MongoDB嵌套文档Upsert更新
在Storm 1.1 + Java 8中实现MongoDB文档的Upsert更新
针对你描述的MongoDB文档结构,要实现timestamp_hour和userid双匹配下的balances字段Upsert,咱们可以通过Storm Bolt结合MongoDB Java驱动来完成。下面是完整的实现方案:
核心思路
- 利用MongoDB的
updateOne方法,搭配upsert: true选项,实现"存在则更新,不存在则插入"的逻辑 - 针对嵌套的
balances字段,使用MongoDB的点路径语法(balances.{key})来精准更新指定键值对 - 在Storm Bolt的
prepare方法初始化MongoDB连接,execute方法处理每个Tuple并执行数据库操作
依赖准备
首先确保你的项目依赖包含以下库(以Maven为例):
<dependencies> <!-- Storm Core 1.1 --> <dependency> <groupId>org.apache.storm</groupId> <artifactId>storm-core</artifactId> <version>1.1.0</version> <scope>provided</scope> </dependency> <!-- MongoDB Java驱动(兼容Storm 1.1的稳定版本) --> <dependency> <groupId>org.mongodb</groupId> <artifactId>mongodb-driver-sync</artifactId> <version>3.12.11</version> </dependency> </dependencies>
完整Bolt实现代码
import org.apache.storm.task.OutputCollector; import org.apache.storm.task.TopologyContext; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.topology.base.BaseRichBolt; import org.apache.storm.tuple.Tuple; import com.mongodb.MongoClient; import com.mongodb.MongoClientURI; import com.mongodb.client.MongoCollection; import com.mongodb.client.MongoDatabase; import com.mongodb.client.model.Filters; import com.mongodb.client.model.UpdateOptions; import com.mongodb.client.model.Updates; import org.bson.Document; import java.util.Map; import java.util.Date; public class MongoBalanceUpsertBolt extends BaseRichBolt { private OutputCollector collector; private MongoCollection<Document> collection; @Override public void prepare(Map<String, Object> topoConf, TopologyContext context, OutputCollector collector) { this.collector = collector; // 初始化MongoDB连接(根据你的配置调整URI) MongoClientURI uri = new MongoClientURI("mongodb://localhost:27017"); MongoClient mongoClient = new MongoClient(uri); MongoDatabase database = mongoClient.getDatabase("your_database_name"); this.collection = database.getCollection("your_collection_name"); } @Override public void execute(Tuple tuple) { try { // 从Tuple中获取业务数据(根据你的Tuple字段名调整) Date timestampHour = (Date) tuple.getValueByField("timestamp_hour"); String userId = tuple.getStringByField("userid"); Integer balanceKey = tuple.getIntegerByField("balance_key"); // 对应balances的键,比如1、2、500 Integer input = tuple.getIntegerByField("input"); Integer output = tuple.getIntegerByField("output"); // 构造查询条件:匹配timestamp_hour和userid var filter = Filters.and( Filters.eq("timestamp_hour", timestampHour), Filters.eq("userid", userId) ); // 构造更新操作:设置balances下指定键的input和output var update = Updates.set( "balances." + balanceKey, new Document("input", input).append("output", output) ); // 执行Upsert操作 collection.updateOne(filter, update, new UpdateOptions().upsert(true)); // 确认Tuple处理成功 collector.ack(tuple); } catch (Exception e) { // 处理异常,fail Tuple以便Storm重试 collector.fail(tuple); e.printStackTrace(); } } @Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // 本Bolt不需要输出,所以留空 } @Override public void cleanup() { // 关闭MongoDB连接(如果需要的话,注意MongoClient是线程安全的,也可以全局复用) if (collection != null && collection.getMongoClient() != null) { collection.getMongoClient().close(); } } }
关键细节说明
- 查询条件:使用
Filters.and组合两个匹配条件,确保只有timestamp_hour和userid完全匹配的文档才会被更新 - 更新语法:
"balances." + balanceKey是MongoDB的嵌套字段路径语法,能精准定位到balances下的指定键,避免覆盖其他已存在的键值对 - Upsert选项:
UpdateOptions().upsert(true)是核心,当没有匹配的文档时,MongoDB会自动插入一条新文档,包含查询条件里的字段和更新的balances内容 - 日期处理:Java的
Date可以直接映射到MongoDB的ISODate类型,无需额外转换 - Storm Tuple处理:务必调用
collector.ack(tuple)和collector.fail(tuple)来保证Storm的可靠性语义
注意事项
- 确保MongoDB的版本和Java驱动版本兼容(这里选的3.12.11是兼容Storm 1.1的稳定版本)
- 如果你的Storm拓扑是分布式部署,建议使用MongoDB连接池或者全局复用MongoClient(MongoClient本身是线程安全的)
- 可以根据业务需求调整Tuple的字段名和数据类型,比如balanceKey如果是字符串类型,直接改对应类型即可
内容的提问来源于stack exchange,提问作者WesternGun
相关产品推荐
相关产品推荐

