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

Java 8+Storm 1.1实现MongoDB嵌套文档Upsert更新

在Storm 1.1 + Java 8中实现MongoDB文档的Upsert更新

针对你描述的MongoDB文档结构,要实现timestamp_hour和userid双匹配下的balances字段Upsert,咱们可以通过Storm Bolt结合MongoDB Java驱动来完成。下面是完整的实现方案:

核心思路

  1. 利用MongoDB的updateOne方法,搭配upsert: true选项,实现"存在则更新,不存在则插入"的逻辑
  2. 针对嵌套的balances字段,使用MongoDB的点路径语法(balances.{key})来精准更新指定键值对
  3. 在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();
        }
    }
}

关键细节说明

  1. 查询条件:使用Filters.and组合两个匹配条件,确保只有timestamp_hour和userid完全匹配的文档才会被更新
  2. 更新语法:"balances." + balanceKey是MongoDB的嵌套字段路径语法,能精准定位到balances下的指定键,避免覆盖其他已存在的键值对
  3. Upsert选项:UpdateOptions().upsert(true)是核心,当没有匹配的文档时,MongoDB会自动插入一条新文档,包含查询条件里的字段和更新的balances内容
  4. 日期处理:Java的Date可以直接映射到MongoDB的ISODate类型,无需额外转换
  5. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:22:59