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

MongoDB Rust驱动Async/Await下插入后聚合查询未等待问题排查

问题:聚合测试因插入未同步导致索引越界失败

我正在为工具模块编写集成测试,该模块的搜索函数返回新聚合管道。当数据库无集合时,测试因accounts[0]索引越界失败,表现为聚合管道未等待插入完成。如何重构让聚合阶段等待插入?

测试代码

#[actix_web::test]
async fn test_search_pipeline() {
    let pipeline = vec![doc! { "$match": { "isDeleted": { "$ne": true } } }];

    let mut updated_pipeline = super::search(&pipeline, String::from("foo"));

    let db: Database = connect_to_database().await;
    let accounts_col: Collection<Document> = db.collection("accounts");

    let account_fixture = account::create();

    // Insert a new account in the database
    let new_account = match accounts_col.insert_one(account_fixture, None).await {
        Ok(new_account) => new_account,
        Err(e) => {
            panic!("Unable to insert new account test fixture document: {}", e);
        }
    };

    // Execute the search aggregate pipeline to find the document we just inserted
    let accounts = match accounts_col.aggregate(updated_pipeline, None).await {
        Ok(mut cursor) => {
            let mut accounts: Vec<Document> = Vec::new();

            while let Some(doc) = cursor.next().await {
                match doc {
                    Ok(doc) => accounts.push(doc),
                    Err(e) => panic!("Error iterating through the new account test pipeline cursor : {}", e)
                }
            }

            accounts
        },
        Err(e) => {
            panic!("Unable to execute test search aggregate pipeline: {}", e);
        }
    };

    // Compare the new account id that we inserted with the first search result account id
    let new_account_id: String = new_account.inserted_id.as_object_id().unwrap().to_hex();
    let search_account_id: String = accounts[0].get_object_id("_id").unwrap().to_hex();

    assert_eq!(new_account_id, search_account_id);
}

测试数据与管道

Account Fixture

doc! {
    "name": "ufc0kpu!pgu1QJW3unj",
    "logo": "60e9ca73c500a9001534ad84-logo.png",
    "isActive": true,
    "isDeleted": false,
    "isUnlimited": false,
    "parentAccount": null,
    "credits": 100,
    "isProAccount": true,
    "is360Enabled": true,
    "isResourcesDisabled": false,
    "isEcoPrinting": false,
    "isCoreExtended": false,
    "createdAt": "2021-09-02T12:11:37.995+0000",
}

生成的聚合管道(updated_pipeline)

doc! {
    "$search": {
        "autocomplete": {
            "query": "ufc0kpu!pgu1QJW3unj",
            "path": "name",
            "fuzzy": {
                "maxEdits": 2,
                "prefixLength": 8
            }
        }
    }
}

已尝试的方案

  • 将聚合阶段嵌套到insert_one的Ok分支,无效
  • 修改insert_one的写入关注为Nodes(1),仍存在竞态条件
  • 尝试因果一致性代码,未保证读已写一致性

解决方案

1. 排查搜索词匹配问题

你的测试代码中传入的搜索词是"foo",但生成的管道中query是fixture的name。如果search函数没有正确将搜索词替换为目标内容,会导致查询不到数据。修正测试代码,使用fixture的name作为搜索词:

let account_fixture = account::create();
// 使用fixture的name作为搜索参数
let search_query = account_fixture.get_str("name").unwrap().to_string();
let mut updated_pipeline = super::search(&pipeline, search_query);

2. 处理全文搜索索引的异步延迟

你使用的$search(自动补全)依赖MongoDB的全文搜索索引,索引构建是异步的——即使插入操作已确认,索引可能还未更新。在测试环境中,可以通过重试查询解决:

// 插入完成后,添加重试逻辑
let mut accounts = Vec::new();
let max_retries = 10;
let mut retry_count = 0;

while accounts.is_empty() && retry_count < max_retries {
    // 等待索引同步
    tokio::time::sleep(tokio::time::Duration::from_millis(500)).await;
    retry_count += 1;

    // 重新执行聚合查询
    accounts = match accounts_col.aggregate(updated_pipeline.clone(), None).await {
        Ok(mut cursor) => {
            let mut accs = Vec::new();
            while let Some(doc) = cursor.next().await {
                accs.push(doc.unwrap());
            }
            accs
        },
        Err(e) => panic!("聚合查询失败: {}", e),
    };
}

// 先断言文档存在,再取索引
assert!(!accounts.is_empty(), "重试10次后仍未找到插入的文档");

// 后续断言逻辑
let new_account_id: String = new_account.inserted_id.as_object_id().unwrap().to_hex();
let search_account_id: String = accounts[0].get_object_id("_id").unwrap().to_hex();
assert_eq!(new_account_id, search_account_id);

3. 先验证文档写入再执行聚合

通过普通find_one确认文档已写入数据库,排除写入本身的问题:

// 插入后先验证文档存在
let inserted_doc = accounts_col.find_one(doc!{"_id": new_account.inserted_id}, None).await.unwrap();
assert!(inserted_doc.is_some(), "插入的文档未成功写入");

// 再执行聚合查询
// ... 原聚合代码

内容的提问来源于stack exchange,提问作者Dylan M

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 03:05:28