Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Fixes CC issue where some relationship inconsistencies may be overlooked
The initial design had a flaw, or a mismatch from one part of the code to another where relationship record processing expected that the records where divided up by its node ids, not its id. The RecordDistributor didn't do this, but the RelationshipRecordCheck code had that expectation. This would have the effect that relationships would randomly not be checked, depending on how many threads the machine had which ran the check.
- Loading branch information
Showing
17 changed files
with
582 additions
and
60 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
113 changes: 113 additions & 0 deletions
113
...onsistency-check/src/main/java/org/neo4j/consistency/checking/full/QueueDistribution.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Original file line | Diff line number | Diff line change |
---|---|---|---|
@@ -0,0 +1,113 @@ | |||
/* | |||
* Copyright (c) 2002-2016 "Neo Technology," | |||
* Network Engine for Objects in Lund AB [http://neotechnology.com] | |||
* | |||
* This file is part of Neo4j. | |||
* | |||
* Neo4j is free software: you can redistribute it and/or modify | |||
* it under the terms of the GNU General Public License as published by | |||
* the Free Software Foundation, either version 3 of the License, or | |||
* (at your option) any later version. | |||
* | |||
* This program is distributed in the hope that it will be useful, | |||
* but WITHOUT ANY WARRANTY; without even the implied warranty of | |||
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the | |||
* GNU General Public License for more details. | |||
* | |||
* You should have received a copy of the GNU General Public License | |||
* along with this program. If not, see <http://www.gnu.org/licenses/>. | |||
*/ | |||
package org.neo4j.consistency.checking.full; | |||
|
|||
import org.neo4j.consistency.checking.full.RecordDistributor.RecordConsumer; | |||
import org.neo4j.kernel.impl.store.record.RelationshipRecord; | |||
|
|||
/** | |||
* Factory for creating {@link QueueDistribution}. Typically the distribution type is decided higher up | |||
* in the call stack and the actual {@link QueueDistributor} is instantiated when more data is available | |||
* deeper down in the call stack. | |||
*/ | |||
public interface QueueDistribution | |||
{ | |||
<RECORD> QueueDistributor<RECORD> distributor( long recordsPerCpu, int numberOfThreads ); | |||
|
|||
/** | |||
* Distributes records into {@link RecordConsumer}. | |||
*/ | |||
public interface QueueDistributor<RECORD> | |||
{ | |||
void distribute( RECORD record, RecordConsumer<RECORD> consumer ) throws InterruptedException; | |||
} | |||
|
|||
/** | |||
* Distributes records round-robin style to all queues. | |||
*/ | |||
public static final QueueDistribution ROUND_ROBIN = new QueueDistribution() | |||
{ | |||
@Override | |||
public <RECORD> QueueDistributor<RECORD> distributor( long recordsPerCpu, int numberOfThreads ) | |||
{ | |||
return new RoundRobinQueueDistributor<>( numberOfThreads ); | |||
} | |||
}; | |||
|
|||
/** | |||
* Distributes {@link RelationshipRecord} depending on the start/end node ids. | |||
*/ | |||
public static final QueueDistribution RELATIONSHIPS = new QueueDistribution() | |||
{ | |||
@Override | |||
public QueueDistributor<RelationshipRecord> distributor( long recordsPerCpu, int numberOfThreads ) | |||
{ | |||
return new RelationshipNodesQueueDistributor( recordsPerCpu ); | |||
} | |||
}; | |||
|
|||
static class RoundRobinQueueDistributor<RECORD> implements QueueDistributor<RECORD> | |||
{ | |||
private final int numberOfThreads; | |||
private int nextQIndex; | |||
|
|||
public RoundRobinQueueDistributor( int numberOfThreads ) | |||
{ | |||
this.numberOfThreads = numberOfThreads; | |||
} | |||
|
|||
@Override | |||
public void distribute( RECORD record, RecordConsumer<RECORD> consumer ) throws InterruptedException | |||
{ | |||
nextQIndex = (nextQIndex + 1) % numberOfThreads; | |||
consumer.accept( record, nextQIndex ); | |||
} | |||
} | |||
|
|||
static class RelationshipNodesQueueDistributor implements QueueDistributor<RelationshipRecord> | |||
{ | |||
private final long recordsPerCpu; | |||
|
|||
public RelationshipNodesQueueDistributor( long recordsPerCpu ) | |||
{ | |||
this.recordsPerCpu = recordsPerCpu; | |||
} | |||
|
|||
@Override | |||
public void distribute( RelationshipRecord relationship, RecordConsumer<RelationshipRecord> consumer ) | |||
throws InterruptedException | |||
{ | |||
int qIndex1 = (int) (relationship.getFirstNode() / recordsPerCpu); | |||
int qIndex2 = (int) (relationship.getSecondNode() / recordsPerCpu); | |||
try | |||
{ | |||
consumer.accept( relationship, qIndex1 ); | |||
if ( qIndex1 != qIndex2 ) | |||
{ | |||
consumer.accept( relationship, qIndex2 ); | |||
} | |||
} | |||
catch ( ArrayIndexOutOfBoundsException e ) | |||
{ | |||
throw e; | |||
} | |||
} | |||
} | |||
} |
Oops, something went wrong.