-
-
Notifications
You must be signed in to change notification settings - Fork 1.2k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Add possibility to specify back pressure class, add benchmark, add mo…
…re tests Signed-off-by: Oleksandr Porunov <alexandr.porunov@gmail.com>
- Loading branch information
Showing
14 changed files
with
746 additions
and
98 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
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1 @@ | ||
[{"unit":"ms/op","name":"org.janusgraph.BackPressureBenchmark.releaseBlocked","value":2754.2091194591017}] |
141 changes: 141 additions & 0 deletions
141
janusgraph-benchmark/src/main/java/org/janusgraph/BackPressureBenchmark.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 | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,141 @@ | ||
// Copyright 2023 JanusGraph Authors | ||
// | ||
// Licensed under the Apache License, Version 2.0 (the "License"); | ||
// you may not use this file except in compliance with the License. | ||
// You may obtain a copy of the License at | ||
// | ||
// http://www.apache.org/licenses/LICENSE-2.0 | ||
// | ||
// Unless required by applicable law or agreed to in writing, software | ||
// distributed under the License is distributed on an "AS IS" BASIS, | ||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
// See the License for the specific language governing permissions and | ||
// limitations under the License. | ||
|
||
package org.janusgraph; | ||
|
||
import org.janusgraph.diskstorage.util.backpressure.SemaphoreQueryBackPressure; | ||
import org.janusgraph.diskstorage.util.backpressure.PassAllQueryBackPressure; | ||
import org.janusgraph.diskstorage.util.backpressure.QueryBackPressure; | ||
import org.janusgraph.diskstorage.util.backpressure.SemaphoreProtectedReleaseQueryBackPressure; | ||
import org.openjdk.jmh.annotations.Benchmark; | ||
import org.openjdk.jmh.annotations.BenchmarkMode; | ||
import org.openjdk.jmh.annotations.Level; | ||
import org.openjdk.jmh.annotations.Mode; | ||
import org.openjdk.jmh.annotations.OutputTimeUnit; | ||
import org.openjdk.jmh.annotations.Param; | ||
import org.openjdk.jmh.annotations.Scope; | ||
import org.openjdk.jmh.annotations.Setup; | ||
import org.openjdk.jmh.annotations.State; | ||
import org.openjdk.jmh.annotations.TearDown; | ||
|
||
import java.util.concurrent.ExecutorService; | ||
import java.util.concurrent.Executors; | ||
import java.util.concurrent.Semaphore; | ||
import java.util.concurrent.TimeUnit; | ||
|
||
import static org.janusgraph.util.system.ExecuteUtil.gracefulExecutorServiceShutdown; | ||
|
||
/** | ||
* Benchmark for different implementations of `QueryBackPressure`. | ||
*/ | ||
@State(Scope.Thread) | ||
@BenchmarkMode(Mode.AverageTime) | ||
@OutputTimeUnit(TimeUnit.MILLISECONDS) | ||
public class BackPressureBenchmark { | ||
|
||
/** | ||
* How many parallel threads which will try to acquire and release queries when the backPressure is reached. | ||
*/ | ||
@Param({ "2000", "1000", "100", "10", "4", "2", "1" }) | ||
int threads; | ||
|
||
/** | ||
* `QueryBackPressure` size (ignored for `passAllBackPressure` type). | ||
*/ | ||
@Param({ "50000", "10000", "1000", "100" }) | ||
int backPressure; | ||
|
||
@Param({ | ||
"semaphoreReleaseProtectedBackPressureWithReleasesAwait", | ||
"semaphoreReleaseProtectedBackPressureWithoutReleasesAwait", | ||
"semaphoreBackPressure", | ||
"passAllBackPressure"}) | ||
String type; | ||
|
||
private QueryBackPressure queryBackPressure; | ||
private ExecutorService queriesAcquireService; | ||
private ExecutorService queriesReleaseService; | ||
private boolean closeBackPressure; | ||
private Semaphore acquireJobsSemaphore; | ||
|
||
@Setup(Level.Invocation) | ||
public void setup() { | ||
acquireJobsSemaphore = new Semaphore(0); | ||
queriesAcquireService = Executors.newFixedThreadPool(threads); | ||
queriesReleaseService = Executors.newFixedThreadPool(threads); | ||
switch (type){ | ||
case "semaphoreReleaseProtectedBackPressureWithReleasesAwait": { | ||
queryBackPressure = new SemaphoreProtectedReleaseQueryBackPressure(backPressure); | ||
closeBackPressure = true; | ||
break; | ||
} | ||
case "semaphoreReleaseProtectedBackPressureWithoutReleasesAwait": { | ||
queryBackPressure = new SemaphoreProtectedReleaseQueryBackPressure(backPressure); | ||
closeBackPressure = false; | ||
break; | ||
} | ||
case "semaphoreBackPressure": { | ||
queryBackPressure = new SemaphoreQueryBackPressure(backPressure); | ||
closeBackPressure = false; | ||
break; | ||
} | ||
case "passAllBackPressure": { | ||
queryBackPressure = new PassAllQueryBackPressure(); | ||
closeBackPressure = false; | ||
break; | ||
} | ||
default: throw new IllegalArgumentException("No implementation found to type = "+type); | ||
} | ||
|
||
for(int j=0; j<backPressure; j++){ | ||
queryBackPressure.acquireBeforeQuery(); | ||
acquireJobsSemaphore.release(); | ||
} | ||
|
||
for(int i = 0; i< threads; i++){ | ||
queriesAcquireService.submit(() -> { | ||
queryBackPressure.acquireBeforeQuery(); | ||
acquireJobsSemaphore.release(); | ||
}); | ||
} | ||
} | ||
|
||
@Benchmark | ||
public void releaseBlocked() { | ||
for(int i = 0; i< threads; i++){ | ||
queriesReleaseService.submit(() -> { | ||
try { | ||
acquireJobsSemaphore.acquire(); | ||
} catch (InterruptedException e) { | ||
throw new RuntimeException(e); | ||
Check warning on line 121 in janusgraph-benchmark/src/main/java/org/janusgraph/BackPressureBenchmark.java
|
||
} | ||
queryBackPressure.releaseAfterQuery(); | ||
}); | ||
} | ||
|
||
gracefulExecutorServiceShutdown(queriesReleaseService, Long.MAX_VALUE); | ||
|
||
if(closeBackPressure){ | ||
queryBackPressure.close(); | ||
} | ||
} | ||
|
||
@TearDown(Level.Invocation) | ||
public void clearResources() { | ||
gracefulExecutorServiceShutdown(queriesAcquireService, Long.MAX_VALUE); | ||
queryBackPressure.close(); | ||
System.gc(); | ||
} | ||
|
||
} |
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
94 changes: 94 additions & 0 deletions
94
.../janusgraph/diskstorage/util/backpressure/SemaphoreProtectedReleaseQueryBackPressure.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 | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,94 @@ | ||
// Copyright 2023 JanusGraph Authors | ||
// | ||
// Licensed under the Apache License, Version 2.0 (the "License"); | ||
// you may not use this file except in compliance with the License. | ||
// You may obtain a copy of the License at | ||
// | ||
// http://www.apache.org/licenses/LICENSE-2.0 | ||
// | ||
// Unless required by applicable law or agreed to in writing, software | ||
// distributed under the License is distributed on an "AS IS" BASIS, | ||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
// See the License for the specific language governing permissions and | ||
// limitations under the License. | ||
|
||
package org.janusgraph.diskstorage.util.backpressure; | ||
|
||
import org.janusgraph.core.JanusGraphException; | ||
import org.slf4j.Logger; | ||
import org.slf4j.LoggerFactory; | ||
|
||
import java.util.concurrent.ExecutorService; | ||
import java.util.concurrent.Executors; | ||
import java.util.concurrent.Semaphore; | ||
|
||
import static org.janusgraph.util.system.ExecuteUtil.gracefulExecutorServiceShutdown; | ||
|
||
/** | ||
* Query back pressure implementation which uses Semaphore to control back pressure and has protection | ||
* in place to not generate more `permits` than `backPressureLimit`.<br> | ||
* | ||
* This implementation is similar to {@link SemaphoreQueryBackPressure } with the exception that `releaseAfterQuery` | ||
* calls are asynchronous (non-blocking) and protected against generating more `permits` than `backPressureLimit`. | ||
* This comes with additional overhead of using a separate thread to process any new `release` calls | ||
* which means that all calls to `releaseAfterQuery` will be processed in sequential order one by one. | ||
* The first time logic registers that an attempt to add a new permit could potentially result in a bigger amount of | ||
* `permits` than `backPressureLimit`, it logs a warning. Subsequent calls to `releaseAfterQuery` will not log such | ||
* warning anymore. | ||
*/ | ||
public class SemaphoreProtectedReleaseQueryBackPressure implements QueryBackPressure{ | ||
|
||
private static final Logger log = LoggerFactory.getLogger(SemaphoreProtectedReleaseQueryBackPressure.class); | ||
|
||
private final ExecutorService executorService = Executors.newSingleThreadExecutor(); | ||
private final Runnable releaseNonBlocking; | ||
private final Semaphore semaphore; | ||
private volatile boolean hadWarningLogged; | ||
|
||
public SemaphoreProtectedReleaseQueryBackPressure(final int backPressureLimit) { | ||
this.semaphore = new Semaphore(backPressureLimit, true); | ||
this.releaseNonBlocking = () -> { | ||
// ensure we never add more permits than `backPressureLimit` | ||
// (even if `releaseAfterQuery()` is called more times than `acquireBeforeQuery()`); | ||
if(semaphore.availablePermits()<backPressureLimit){ | ||
semaphore.release(); | ||
} else if(!hadWarningLogged){ | ||
log.warn("`releaseAfterQuery` is called more than once for some of the `acquireBeforeQuery` calls. " + | ||
"This is a sign that the logic using this `QueryBackPressure` may not properly handle special " + | ||
"(potentially exceptional) cases. {} will not trigger more releases than {}. This warning will " + | ||
"be logged only once and it will be ignored for other `releaseAfterQuery` calls which attempt " + | ||
"to add more permits than the configured limit.", | ||
SemaphoreProtectedReleaseQueryBackPressure.class.getSimpleName(), backPressureLimit); | ||
hadWarningLogged = true; | ||
} | ||
}; | ||
} | ||
|
||
@Override | ||
public void acquireBeforeQuery() { | ||
try { | ||
semaphore.acquire(); | ||
} catch (InterruptedException e) { | ||
Thread.currentThread().interrupt(); | ||
throw new JanusGraphException(e); | ||
} | ||
} | ||
|
||
@Override | ||
public void releaseAfterQuery(){ | ||
executorService.execute(releaseNonBlocking); | ||
} | ||
|
||
@Override | ||
public void close() { | ||
gracefulExecutorServiceShutdown(executorService, Long.MAX_VALUE); | ||
} | ||
|
||
int availablePermits(){ | ||
return semaphore.availablePermits(); | ||
} | ||
|
||
boolean hadWarningLogged(){ | ||
return hadWarningLogged; | ||
} | ||
} |
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
Oops, something went wrong.