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

并行创建一对多Edge时触发ConcurrentModificationException问题求助

并发创建Edge时触发VertexProperty并发修改异常

我在从指定顶点并发创建数百条指向不同唯一顶点的Edge时,遇到了如下错误:

{
"requestId": "b90671d4-d8f0-4cca-b0c3-c908ff91d022",
"code": "ConcurrentModificationException",
"detailedMessage": "Failed to complete Insert operation for a VertexProperty due to conflicting concurrent operations. Please retry. 0 transactions are currently rolling back.",
"message": "Failed to complete Insert operation for a VertexProperty due to conflicting concurrent operations. Please retry. 0 transactions are currently rolling back."
}

我确认自己仅在创建Edge属性,没有重复修改同一条Edge,也没有同时写入同一顶点对的多条Edge。怀疑是错误信息描述有误,或者Gremlin在同一顶点创建多Edge并写入属性时存在问题;也可能我误解了查询逻辑,误操作了VertexProperty而非EdgeProperty,但测试相似查询后排除了这个可能。

以下是我正在调试的代码:

router.post(
  '/users/contacts',
  asyncHandler(async (req: Request, res: Response) => {
    var successfulPhoneNumbers: string[] = [];
    var failedPhoneNumbers: string[] = [];
    try {
      const g = req.g;

      let { userId, contacts }: { userId: string; contacts: Contact[] } = req.body;

      if (!userId || !contacts) {
        return res.status(400).end('Required userId or contacts not supplied');
      }

      if (contacts.length > 300) {
        return res.status(400).end('Maximum 300 contacts allowed');
      }

      let limit = pLimit(300);

      var now = new Date(Date.now()).toISOString();

      // filter out duplicates
      contacts = contacts.filter(
        (contact, index, self) =>
          index === self.findIndex((t) => t.phoneNumber === contact.phoneNumber),
      );

      let contactsPromises = contacts.map(async (contact) => {
        let addContactProperties = (
          query: gremlin.process.GraphTraversal,
        ): gremlin.process.GraphTraversal => {
          query = query.property('createdAt', now);

          if (contact.firstName) {
            query = query.property('firstName', contact.firstName);
          }
          if (contact.lastName) {
            query = query.property('lastName', contact.lastName);
          }
          if (contact.email) {
            query = query.property('email', contact.email);
          }
          if (contact.birthday) {
            query = query.property('birthday', contact.birthday);
          }
          if (contact.address) {
            query = query.property('address', contact.address);
          }

          return query;
        };

        return limit(async () => {
          try {
            let phoneNumber = contact.phoneNumber;
            if (!isPhoneNumberValid(phoneNumber)) {
              failedPhoneNumbers.push(phoneNumber);
              return;
            }

            // if the contact is on the system, add an edge to their user node

            var contactAsUser = await g.V().hasLabel('user').has('phoneNumber', phoneNumber).next();
            var contactAsContact = await g
              .V()
              .hasLabel('contact')
              .has('phoneNumber', phoneNumber)
              .next();

            if (!contactAsUser.done) {
              let query = g
                .V(userId)
                .coalesce(
                  __.outE('HAS_CONTACT').where(__.inV().hasId(contactAsUser.value.id)),
                  __.addE('HAS_CONTACT').to(__.V(contactAsUser.value.id)),
                );
              query = addContactProperties(query);
              await query.iterate();
            } else if (!contactAsContact.done) {
              let query = g
                .V(userId)
                .coalesce(
                  __.outE('HAS_CONTACT').where(__.inV().hasId(contactAsContact.value.id)),
                  __.addE('HAS_CONTACT').to(__.V(contactAsContact.value.id)),
                );
              query = addContactProperties(query);
              await query.iterate();
            } else {
              // if the contact is not on the system, add an edge to a new contact node
              let query = g
                .V(userId)
                .coalesce(
                  __.outE('HAS_CONTACT').where(__.inV().has('phoneNumber', phoneNumber)),
                  __.addV('contact')
                    .as('contact')
                    .property(single, 'phoneNumber', phoneNumber)
                    .addE('HAS_CONTACT')
                    .from_(__.V(userId))
                    .to(__.select('contact')),
                );
              query = addContactProperties(query);
              await query.iterate();
            }
            successfulPhoneNumbers.push(phoneNumber);
          } catch (err: any) {
            console.log(err.toString());
            failedPhoneNumbers.push(contact.phoneNumber);
          }
        });
      });

      await Promise.all(contactsPromises);

      return res.json({
        successfulPhoneNumbers: successfulPhoneNumbers,
        failedPhoneNumbers: failedPhoneNumbers,
      });
    } catch (err: any) {
      return res.status(500).json({
        error: err.message ?? err,
        successfulPhoneNumbers: successfulPhoneNumbers,
        failedPhoneNumbers: failedPhoneNumbers,
      });
    } finally {
      await req.dc.close();
    }
  }),
);

问题根源

你遇到的异常本质是并发修改了同一个起始顶点(userId对应的顶点)的内部结构。虽然你看起来只是在给不同的Edge加属性,但Gremlin在处理从同一个顶点出发的Edge创建/更新时,会修改该顶点的邻接表(属于顶点的内部元数据或隐含属性)。当数百个并发请求同时操作这个顶点时,就会触发VertexProperty的并发冲突——这里的VertexProperty并非你显式设置的属性,而是图数据库底层维护的顶点关联Edge的内部结构。

另外,代码里的几个点放大了冲突概率:

  • 直接用pLimit(300)并发300个请求,所有请求都操作同一个起始顶点,竞争过于激烈
  • 每个请求先执行两次独立查询(检查contact是否为user/contact),再执行Edge操作,增加了事务冲突的窗口
  • 第三个分支(创建新contact顶点+Edge)的coalesce判断存在竞态:多个并发请求可能同时判断到该phoneNumber不存在,进而尝试创建同一个contact顶点,触发额外冲突

解决方案

1. 改用批量单查询处理(推荐)

将批量操作合并为单个Gremlin查询,让数据库端原子性完成所有操作,彻底避免客户端并发冲突。示例思路:

// 替换原来的contactsPromises和Promise.all部分
const validContacts = contacts.filter(c => isPhoneNumberValid(c.phoneNumber));

await g.V(userId)
  .inject(validContacts)
  .unfold()
  .as('contactData')
  // 原子化检查目标顶点是否存在,不存在则创建
  .choose(
    __.V().has('phoneNumber', __.select('contactData').values('phoneNumber')).fold(),
    __.unfold().as('target'),
    __.addV('contact').property(single, 'phoneNumber', __.select('contactData').values('phoneNumber')).as('target')
  )
  // 原子化检查边是否存在,不存在则创建
  .coalesce(
    __.V(userId).outE('HAS_CONTACT').where(__.inV().as('target')),
    __.V(userId).addE('HAS_CONTACT').to(__.select('target'))
  )
  // 统一设置Edge属性,用coalesce避免空值写入
  .property('createdAt', now)
  .coalesce(__.values('firstName'), __.property('firstName', __.select('contactData').values('firstName')))
  .coalesce(__.values('lastName'), __.property('lastName', __.select('contactData').values('lastName')))
  .coalesce(__.values('email'), __.property('email', __.select('contactData').values('email')))
  .coalesce(__.values('birthday'), __.property('birthday', __.select('contactData').values('birthday')))
  .coalesce(__.values('address'), __.property('address', __.select('contactData').values('address')))
  .iterate();

// 后续区分成功失败:有效号码默认成功,可额外查询验证
successfulPhoneNumbers = validContacts.map(c => c.phoneNumber);
failedPhoneNumbers = contacts.filter(c => !isPhoneNumberValid(c.phoneNumber)).map(c => c.phoneNumber);

2. 优化现有并发逻辑(保留并发场景)

如果必须保留并发模式,大幅降低并发数并增加重试机制:

// 增加重试函数,针对并发异常自动重试
const retryOnConflict = async (operation, retries = 3) => {
  try {
    return await operation();
  } catch (err) {
    if (retries > 0 && err.code === 'ConcurrentModificationException') {
      await new Promise(resolve => setTimeout(resolve, 100 * (4 - retries)));
      return retryOnConflict(operation, retries - 1);
    }
    throw err;
  }
};

// 把并发数从300降到10-20,减少顶点竞争
let limit = pLimit(15);

// 在limit回调中使用重试
return limit(async () => {
  try {
    await retryOnConflict(async () => {
      // 原有的操作逻辑(检查user/contact、创建Edge等)
    });
    successfulPhoneNumbers.push(phoneNumber);
  } catch (err) {
    failedPhoneNumbers.push(contact.phoneNumber);
  }
});

3. 修复竞态条件

第三个分支的coalesce判断需要原子化,避免多个请求同时创建同一contact顶点:

// 替换原第三个分支的query逻辑
let query = g
  .V(userId)
  .coalesce(
    __.outE('HAS_CONTACT').where(__.inV().has('phoneNumber', phoneNumber)),
    // 原子化检查+创建contact顶点
    __.V().has('phoneNumber', phoneNumber).fold().coalesce(
      __.unfold(),
      __.addV('contact').property(single, 'phoneNumber', phoneNumber)
    ).as('target')
    .addE('HAS_CONTACT').from_(__.V(userId)).to(__.select('target'))
  );
query = addContactProperties(query);
await query.iterate();

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 21:22:48