并行创建一对多Edge时触发ConcurrentModificationException问题求助
我在从指定顶点并发创建数百条指向不同唯一顶点的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

