Neo4j PARALLEL运行时结果顺序异常问题排查与解决
Neo4j PARALLEL运行时查询结果不稳定问题解决
我在测试Neo4j的PARALLEL运行时,但测试无法通过——查询返回的结果顺序和SLOTTED/PIPELINED运行时不一致,甚至weight值也不稳定。比如多次执行下面的CYPHER查询,有时候weight=25的记录排在首位,有时候却是weight=10的记录。换成CYPHER runtime = pipelined时查询完全正常。
CYPHER runtime = parallel MATCH (dg:DecisionGroup { id: -2 })-[rdgd: CONTAINS ]-> ( childD:Profile ) MATCH (childD)-[mhvo:HAS_VOTE_ON]-(mc:Criterion) WHERE mc.id IN [9760, 9761, 9757, 9758, 9759] WITH childD , collect(mhvo) AS mhvos WHERE size(mhvos) >= size([9760, 9761, 9757, 9758, 9759]) WITH childD WHERE (childD.`active` = true ) MATCH (childDStat:JobableStatistic { jobableId: childD.id }) WITH childD, childDStat UNWIND [9760, 9761, 9757, 9758, 9759] AS dCId WITH childD, childDStat, dCId + coalesce({}[toString(dCId)], []) AS cGroup WHERE NOT AlL(x IN cGroup WHERE x IN childDStat.zeroCriterionIds ) WITH childD, childDStat, collect(cGroup) AS cGroups WHERE size(cGroups) >= size([9760, 9761, 9757, 9758, 9759]) UNWIND cGroups AS cGroup WITH childD, childDStat, cGroup WHERE ANY(x IN cGroup WHERE x IN childDStat.detailedCriterionIds) WITH childD, childDStat, collect(cGroup) AS cGroups WHERE size(cGroups) >= size([9760, 9761, 9757, 9758, 9759]) WITH childD, childDStat, size(cGroups) AS cGroupsSize, cGroups UNWIND cGroups AS cGroup WITH childD, childDStat, cGroups, cGroupsSize, cGroup UNWIND cGroup AS cId WITH childD, childDStat, cGroups, cGroupsSize, cGroup, cId, cGroup[0] AS cG0 WITH childD, childDStat, cGroups, cGroupsSize, cGroup, cId, cG0, { `9759`:1.0, `9758`:1.0, `9760`:1.0, `9761`:1.0, `9757`:1.0 }[toString(cG0)] AS criterionAvgVoteWeight, { `9759`:0, `9758`:0, `9760`:0, `9761`:0, `9757`:0 }[toString(cG0)] AS criterionExperienceMonth WHERE (criterionAvgVoteWeight = 0 OR criterionAvgVoteWeight <= childDStat['criterionAvgVoteWeights.' + cId]) AND (criterionExperienceMonth = 0 OR criterionExperienceMonth <= childDStat['criterionExperienceMonths.' + cId]) WITH childD, childDStat, cGroups, cGroupsSize, collect(cId) AS cIds UNWIND cGroups AS cGroup WITH childD, childDStat, cGroup, cGroupsSize, cIds WHERE ANY(x IN cIds WHERE x IN cGroup) WITH childD, childDStat, cGroupsSize, collect(cGroup) AS cGroups WHERE size(cGroups) >= cGroupsSize WITH childD, childDStat UNWIND [9760, 9761, 9757, 9758, 9759] AS cId WITH childD, childDStat, cId, { `9759`:1.0, `9758`:1.0, `9760`:1.0, `9761`:1.0, `9757`:1.0 }[toString(cId)] AS criterionCoefficient WITH childD, sum(criterionCoefficient * childDStat['criterionAvgVoteWeights.' + cId]) AS weight, sum(childDStat['criterionTotalVotes.' + cId]) AS totalVotes, sum(criterionCoefficient) AS criterionCoefficientSum WITH childD, weight, totalVotes, criterionCoefficientSum WITH collect({`childD`:childD , `weight`:weight, `totalVotes`: totalVotes }) AS aggregate WITH aggregate, size(aggregate) AS count UNWIND aggregate AS item WITH count, item.childD AS childD , item.weight AS weight, item.totalVotes AS totalVotes MATCH (dg:DecisionGroup { id: -2 })-[rdgd: CONTAINS ]->(childD) OPTIONAL MATCH (childD)-[ru:CREATED_BY]->(u:User) OPTIONAL MATCH (jobable:Decision:Vacancy { id: 4928 }) RETURN count, childD AS decision, dg, rdgd , u, ru , jobable.id AS jobableId , weight, totalVotes, [ (jobable)-[vg1:HAS_VOTE_ON]->(c1:Criterion) | { criterion: c1, relationship: vg1 } ] AS jobableWeightedCriteria, [(jobable)-[:HAS_VOTE_ON]->(c1:Criterion)<-[vg1:HAS_VOTE_ON]-(childD) WHERE c1.id IN childD.detailedCriterionIds | { criterion: c1, relationship: vg1 } ] AS weightedCriteria , [ (c1t:Translation:BaseEntity)<-[rc1t: CONTAINS ]-(c1:Criterion)<-[vg1:HAS_VOTE_ON]-(childD) WHERE EXISTS ((jobable)-[:HAS_VOTE_ON]->(c1)) AND c1t.iso6391 = 'uk' AND c1.id IN childD.detailedCriterionIds | { entityId: toInteger(c1.id), translation: c1t } ] AS weightedCriteriaTranslations , [ (jobable)-[:WORK_LOCATED_IN|EMPLOYMENT_TYPE_AS|READY_TO|EMPLOYMENT_AS|WORK_TIME_ZONE|BELONGS_TO|LOCATED_IN|COMPANY|WORK_PERMIT_IN|COMPANY_TYPE_OF]->(ce:CompositeEntity) | { entity: ce } ] AS jobableCompositeEntities, [ (childD)-[:WORK_LOCATED_IN|EMPLOYMENT_TYPE_AS|READY_TO|EMPLOYMENT_AS|WORK_TIME_ZONE|BELONGS_TO|LOCATED_IN|COMPANY|WORK_PERMIT_IN|COMPANY_TYPE_OF]->(ce:CompositeEntity) | { entity: ce } ] AS decisionCompositeEntities, [ (childD)-[:WORK_LOCATED_IN|EMPLOYMENT_TYPE_AS|READY_TO|EMPLOYMENT_AS|WORK_TIME_ZONE|BELONGS_TO|LOCATED_IN|COMPANY|WORK_PERMIT_IN|COMPANY_TYPE_OF]->(ce:CompositeEntity)-[: CONTAINS ]->(trans:Translation:BaseEntity) WHERE trans.iso6391 = 'uk' | { entityId: toInteger(id(ce)), translation: trans } ] AS decisionCompositeEntitiesTranslations, [ (childD)-[: CONTAINS ]->(trans:Translation:BaseEntity) WHERE trans.iso6391 = 'uk' | { entityId: toInteger(childD.id), translation: trans } ] AS decisionTranslations, [ (rc:Criterion)-[*0]->() WHERE rc.id IN childD.replaceableCriterionIds | { entity: rc } ] AS decisionReplaceableCriteria, [ (rc:Criterion)-[: CONTAINS ]->(trans:Translation:BaseEntity) WHERE rc.id IN childD.replaceableCriterionIds AND trans.iso6391 = 'uk' | { entityId: toInteger(id(rc)), translation: trans } ] AS decisionReplaceableCriteriaTranslations, COUNT { (:Vacancy:Jobable:BaseEntity {status: 'APPROVED', active: true })<-[:POTENTIAL_PROFILE]-(childD) } AS potentialJobablesCount , COUNT { (:Vacancy:Jobable:BaseEntity {status: 'APPROVED', active: true })<-[:RELEVANT_PROFILE]-(childD) } AS relevantJobablesCount ORDER BY weight DESC, childD.createdAt DESC SKIP 0 LIMIT 100
我怀疑问题出在模式组合或者COUNT子查询上,尤其是RETURN语句里的这些部分:
, [ (jobable)-[vg1:HAS_VOTE_ON]->(c1:Criterion) | { criterion: c1, relationship: vg1 } ] AS jobableWeightedCriteria, [(jobable)-[:HAS_VOTE_ON]->(c1:Criterion)<-[vg1:HAS_VOTE_ON]-(childD) WHERE c1.id IN childD.detailedCriterionIds | { criterion: c1, relationship: vg1 } ] AS weightedCriteria , [ (c1t:Translation:BaseEntity)<-[rc1t: CONTAINS ]-(c1:Criterion)<-[vg1:HAS_VOTE_ON]-(childD) WHERE EXISTS ((jobable)-[:HAS_VOTE_ON]->(c1)) AND c1t.iso6391 = 'uk' AND c1.id IN childD.detailedCriterionIds | { entityId: toInteger(c1.id), translation: c1t } ] AS weightedCriteriaTranslations , [ (jobable)-[:WORK_LOCATED_IN|EMPLOYMENT_TYPE_AS|READY_TO|EMPLOYMENT_AS|WORK_TIME_ZONE|BELONGS_TO|LOCATED_IN|COMPANY|WORK_PERMIT_IN|COMPANY_TYPE_OF]->(ce:CompositeEntity) | { entity: ce } ] AS jobableCompositeEntities, [ (childD)-[:WORK_LOCATED_IN|EMPLOYMENT_TYPE_AS|READY_TO|EMPLOYMENT_AS|WORK_TIME_ZONE|BELONGS_TO|LOCATED_IN|COMPANY|WORK_PERMIT_IN|COMPANY_TYPE_OF]->(ce:CompositeEntity) | { entity: ce } ] AS decisionCompositeEntities, [ (childD)-[:WORK_LOCATED_IN|EMPLOYMENT_TYPE_AS|READY_TO|EMPLOYMENT_AS|WORK_TIME_ZONE|BELONGS_TO|LOCATED_IN|COMPANY|WORK_PERMIT_IN|COMPANY_TYPE_OF]->(ce:CompositeEntity)-[: CONTAINS ]->(trans:Translation:BaseEntity) WHERE trans.iso6391 = 'uk' | { entityId: toInteger(id(ce)), translation: trans } ] AS decisionCompositeEntitiesTranslations, [ (childD)-[: CONTAINS ]->(trans:Translation:BaseEntity) WHERE trans.iso6391 = 'uk' | { entityId: toInteger(childD.id), translation: trans } ] AS decisionTranslations, [ (rc:Criterion)-[*0]->() WHERE rc.id IN childD.replaceableCriterionIds | { entity: rc } ] AS decisionReplaceableCriteria, [ (rc:Criterion)-[: CONTAINS ]->(trans:Translation:BaseEntity) WHERE rc.id IN childD.replaceableCriterionIds AND trans.iso6391 = 'uk' | { entityId: toInteger(id(rc)), translation: trans } ] AS decisionReplaceableCriteriaTranslations, COUNT { (:Vacancy:Jobable:BaseEntity {status: 'APPROVED', active: true })<-[:POTENTIAL_PROFILE]-(childD) } AS potentialJobablesCount , COUNT { (:Vacancy:Jobable:BaseEntity {status: 'APPROVED', active: true })<-[:RELEVANT_PROFILE]-(childD) } AS relevantJobablesCount
解决方法
1. 强制稳定排序
当前的排序条件ORDER BY weight DESC, childD.createdAt DESC可能存在重复值,PARALLEL运行时并行处理时,排序键相同的记录顺序会随机。添加唯一标识字段(比如childD.id)到排序末尾,确保排序结果唯一稳定:
ORDER BY weight DESC, childD.createdAt DESC, childD.id DESC
2. 提前计算聚合与复杂逻辑
把RETURN阶段的COUNT子查询、集合生成逻辑提前到WITH阶段处理,减少RETURN阶段并行计算的冲突:
- 提前计算两个COUNT值:
WITH childD, weight, totalVotes, COUNT { (:Vacancy:Jobable:BaseEntity {status: 'APPROVED', active: true})<-[:POTENTIAL_PROFILE]-(childD) } AS potentialJobablesCount, COUNT { (:Vacancy:Jobable:BaseEntity {status: 'APPROVED', active: true})<-[:RELEVANT_PROFILE]-(childD) } AS relevantJobablesCount
- 类似地,将
weightedCriteria等集合生成逻辑也提前到WITH中,避免在RETURN阶段并行生成时干扰排序。
3. 简化查询冗余部分
删除原查询中重复的WITH childD语句;将[9760, 9761, 9757, 9758, 9759] + []简化为[9760, 9761, 9757, 9758, 9759],减少不必要的计算,降低并行处理的不确定性。
4. 验证版本兼容性
确保使用的Neo4j版本(建议4.4+)对PARALLEL运行时的支持成熟,旧版本可能存在并行聚合或排序的bug。
内容的提问来源于stack exchange,提问作者alexanoid
相关产品推荐
相关产品推荐

