【问题标题】:Query oplog timestamp with spring mongo使用spring mongo查询oplog时间戳
【发布时间】:2017-08-07 08:15:31
【问题描述】:

我想使用 Java 查询 mongodbs oplog,如果可能的话,使用 spring mongodb 集成。我的问题是从 java 创建以下查询:

db['oplog.rs'].find({ "ts": { $gt: Timestamp(1489568405,34) }, $and: [ { "ns": "myns" } ] })

我尝试了一些方法,例如 BsonTimestamp 或 BSONTimestamp,它们会导致错误的查询。使用

BasicQuery({ "ts": { $gt: Timestamp(1489568405,34) }, $and: [ { "ns": "myns" } ] }) 

导致java mongodb驱动的JSON解析器出错。

有什么提示吗?

感谢于尔根

典型的记录如下所示:

{ 
    "ts" : Timestamp(1489567144, 2), 
    "t" : NumberLong(2), 
    "h" : NumberLong(7303473893196954969), 
    "v" : NumberInt(2), 
    "op" : "i", 
    "ns" : "asda.jam", 
    "o" : {
        "_id" : NumberInt(2), 
        "time" : ISODate("2017-03-15T08:39:00.000+0000"), 
        "roadDesc" : {
            "roadId" : NumberInt(28102917), 
            "roadName" : "A480 W"
        }, 
        "posUpFront" : NumberInt(1003), 
        "posDownFront" : NumberInt(1003), 
        "_class" : "de.heuboe.acaJNI.test.Jam"
    }
}

【问题讨论】:

    标签: java spring mongodb timestamp spring-data-mongodb


    【解决方案1】:

    Mongo 为 NumberLong、Timestamp 等结构提供了扩展的 JSON 语法,可在 Mongo shell 上运行。为了使其在 Java 代码中工作,它们具有严格的 JSON 模式,其中这些运算符使用 JSON (https://docs.mongodb.com/manual/reference/mongodb-extended-json/#bson-data-types-and-associated-representations) 表示。要使用 Java 执行此操作,您可以创建一个自定义转换器并将其注册到您的 MappingMongoConverter(请参阅下面的 sn-p)。转换器应将数据类型(例如 BSONTimestamp)转换为适当的严格 JSON 文档格式。

    @WritingConverter
    public class BsonTimestampToDocumentConverter implements Converter<BSONTimestamp, Document> {
    
      private static final Logger LOGGER = LoggerFactory.getLogger(BsonTimestampToDocumentConverter.class);
    
      public BsonTimestampToDocumentConverter() {
        //
      }
    
      @Override
      public Document convert(BSONTimestamp source) {
        LOGGER.trace(">>>> Converting BSONTimestamp to Document");
        Document value = new Document();
        value.put("t", source.getTime());
        value.put("i", source.getInc());
        return new Document("$timestamp", value);
      }
    }
    

    像这样在 MappingMongoConverter 中注册

     public MappingMongoConverter syncLocalMappingMongoConverter() throws Exception {
        MongoMappingContext mappingContext = new MongoMappingContext();
        DbRefResolver dbRefResolver = new DefaultDbRefResolver(syncLocalDbFactory());
        MappingMongoConverter converter = new MappingMongoConverter(dbRefResolver, mappingContext);
        converter.setCustomConversions(customConversions());
    
        return converter;
     }
    
    
      private CustomConversions customConversions()  {
       List<Converter<?, ?>> converterList = new ArrayList<>();
       converterList.add(new BsonTimestampToDocumentConverter());
       // add the other converters here
       return new CustomConversions(CustomConversions.StoreConversions.NONE, converterList);
    }
    

    这是一个示例,我使用它查询 oplog 存储库以在特定时间后返回记录(存储库中的 Sync 用于将其与我正在处理的反应性异步内容区分开来。Async 存储库看起来完全相同,除了它应该扩展 ReactiveMongoRepository)。 OplogRecord 类是我创建的一个 Java bean,用于匹配 MongoDb oplog 记录的结构。

    public interface SyncOplogRepository extends MongoRepository<OplogRecord, Long> {
    
      @Query(value = "{ \"op\": { $nin: ['n', 'c'] } }") List<OplogRecord> findRecordsNotEqualToNOrC();
    
      @Query(value = "{'ts' : {$gte : ?0}, \"op\": { $nin: ['n', 'c'] } }")
      List<OplogRecord> findRecordsNotEqualToNOrCAfterTime(BSONTimestamp timestamp);
    
      @Query(value = "{'ts' : {$lt : ?0}, \"op\": { $nin: ['n', 'c'] } }") 
      List<OplogRecord> findRecordsNotEqualToNOrCBeforeTime(BSONTimestamp timestamp);
    
    }
    

    OplogRecord 类

    import com.mongodb.DBObject;
    import org.bson.BsonTimestamp;
    import org.bson.types.BSONTimestamp;
    import org.slf4j.Logger;
    import org.slf4j.LoggerFactory;
    import org.springframework.data.annotation.Id;
    import org.springframework.data.mongodb.core.mapping.Document;
    
    import java.util.Map;
    
    
    @Document(collection = "oplog.rs")
    public class OplogRecord {
    
      @Id
      private Long id;
    
      /**
       * Timestamp
       */
      private BsonTimestamp ts;
    
      /**
       * Unique id for this entry
       */
      private Long h;
    
      /**
       * DB and collection name of change.
       */
      private String ns;
    
      /**
       * The actual document that was modified/inserted/deleted
       */
      private Map<String, Object> o;
    
      /**
       * The operation that was performed
       */
      private String op;
    
      /**
       * ??
       */
      private Long t;
    
      /**
       * ??
       */
      private Integer v;
    
      public BsonTimestamp getTs() {
        return ts;
      }
    
      public void setTs(BsonTimestamp ts) {
        this.ts = ts;
      }
    
      public Long getH() {
        return h;
      }
    
      public void setH(Long h) {
        this.h = h;
      }
    
      public String getNs() {
        return ns;
      }
    
      public void setNs(String ns) {
        this.ns = ns;
      }
    
      public Map<String, Object> getO() {
        return o;
      }
    
      public void setO(Map<String, Object> o) {
        this.o = o;
      }
    
      public String getOp() {
        return op;
      }
    
      public void setOp(String op) {
        this.op = op;
      }
    
      public Long getT() {
        return t;
      }
    
      public void setT(Long t) {
        this.t = t;
      }
    
      public Integer getV() {
        return v;
      }
    
      public void setV(Integer v) {
        this.v = v;
      }
    }
    

    ~
    ~

    【讨论】:

    • 感谢您的回答。我尝试使用 CustomConverter 但没有成功。你能告诉我你是如何建模 OplogRecord 的吗?
    • 添加了 OplogRecord 类。这只是一个 POJO。确保使用的 MappingMongoContext 实际具有您注册的转换器。
    • 你能详细说明吗?
    • 对不起,我昨天走神了。仍然无法让它工作。我的查询结果为:{ "ts" : { "$gte" : { "$timestamp" : { "t" : 1490164923 , "i" : 32}}}},导致结果为空。在此之前我是否必须将驱动程序或数据库切换到严格的 json 模式?
    【解决方案2】:

    您可以使用 org.bson.BsonTimestamp 进行过滤。

    BsonTimestamp lastReadTimestamp = new BsonTimestamp(1489568405, 34);
    Bson filter = new Document("$gt", lastReadTimestamp);
    

    然后你可以像这样使用 find,

    oplogColl.find(new Document("ts", filter));
    

    或者您可以创建一个可保释光标并像这样遍历文档,

    MongoCursor oplogCursor =
                        oplogColl
                                .find(new Document("ts", filter))
                                .cursorType(CursorType.TailableAwait)
                                .noCursorTimeout(true)
                                .batchSize(1000)
                                .iterator();
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-06-21
      • 2018-12-23
      • 2014-10-06
      • 1970-01-01
      • 2016-08-01
      • 2018-04-21
      • 2021-12-11
      • 2015-11-18
      相关资源
      最近更新 更多