Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion tests/docker-images/java-test-image/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ USER root
COPY target/scripts /pulsar/bin
RUN chmod a+rx /pulsar/bin/*

RUN apk add --no-cache supervisor
RUN apk add --no-cache supervisor jq

RUN mkdir -p /var/log/pulsar \
&& mkdir -p /var/run/supervisor/ \
Expand Down
2 changes: 1 addition & 1 deletion tests/docker-images/latest-version-image/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,7 @@ FROM $PULSAR_ALL_IMAGE
# However, any processes exec'ing into the containers will run as root, by default.
USER root

RUN apk add --no-cache supervisor procps curl
RUN apk add --no-cache supervisor procps curl jq

RUN mkdir -p /var/log/pulsar && mkdir -p /var/run/supervisor/

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ loglevel=info
pidfile=/var/run/supervisord.pid
minfds=1024
minprocs=200
user=root

[unix_http_server]
file=/var/run/supervisor/supervisor.sock
Expand Down
19 changes: 14 additions & 5 deletions tests/docker-images/latest-version-image/scripts/func-lib.sh
Original file line number Diff line number Diff line change
Expand Up @@ -21,21 +21,30 @@
set -e
set -o pipefail

function set_pulsar_mem() {
local maxMem=$1
local additionalMemParam=$2
local pulsar_test_mem
# set into pulsar_test_mem while trimming whitespace
read -r pulsar_test_mem <<< "-Xmx${maxMem} ${additionalMemParam}"
# prefer PULSAR_MEM, but always append params to perform a heap dump on OOME
export PULSAR_MEM="${PULSAR_MEM:-"${pulsar_test_mem}"} -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/var/log/pulsar -XX:+ExitOnOutOfMemoryError"
}

function run_pulsar_component() {
local component=$1
local supervisord_component=$2
local maxMem=$3
local additionalMemParam=$4
export PULSAR_MEM="${PULSAR_MEM:-"-Xmx${maxMem} ${additionalMemParam} -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/var/log/pulsar -XX:+ExitOnOutOfMemoryError"}"
export PULSAR_GC="${PULSAR_GC:-"-XX:+UseZGC"}"

set_pulsar_mem "$maxMem" "$additionalMemParam"

if [[ -f "conf/${component}.conf" ]]; then
bin/apply-config-from-env.py conf/${component}.conf
fi
bin/apply-config-from-env.py conf/pulsar_env.sh
bin/apply-config-from-env.py conf/client.conf

if [[ "$component" == "functions_worker" ]]; then
bin/apply-config-from-env.py conf/client.conf
bin/gen-yml-from-env.py conf/functions_worker.yml
fi

Expand All @@ -46,7 +55,7 @@ function run_pulsar_component() {
fi

if [ -z "$NO_AUTOSTART" ]; then
sed -i 's/autostart=.*/autostart=true/' /etc/supervisord/conf.d/${supervisord_component}.conf
sed -i 's/autostart=.*/autostart=true/' /etc/supervisord/conf.d/${supervisord_component}.conf
fi

exec /usr/bin/supervisord -c /etc/supervisord.conf
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,8 @@
# under the License.
#

export PULSAR_MEM="${PULSAR_MEM:-"-Xmx512M -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/var/log/pulsar -XX:+ExitOnOutOfMemoryError"}"
export PULSAR_GC="${PULSAR_GC:-"-XX:+UseZGC"}"
source /pulsar/bin/func-lib.sh

set_pulsar_mem 512M

bin/pulsar standalone
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,6 @@ protected ChaosContainer(String clusterName, String image) {
@Override
protected void configure() {
super.configure();
addEnv("MALLOC_ARENA_MAX", "1");
}

protected void appendToEnv(String key, String value) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,11 @@
*/
package org.apache.pulsar.tests.integration.profiling;

import java.io.File;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.attribute.PosixFilePermissions;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
Expand All @@ -31,6 +36,7 @@
import org.apache.pulsar.tests.integration.suites.PulsarTestSuite;
import org.apache.pulsar.tests.integration.topologies.PulsarClusterSpec;
import org.apache.pulsar.tests.integration.utils.DockerUtils;
import org.testcontainers.containers.BindMode;
import org.testcontainers.containers.GenericContainer;
import org.testng.annotations.Test;

Expand All @@ -40,9 +46,9 @@
* Example usage:
* # This has been tested on Mac with Orbstack (https://orbstack.dev/) docker
* # compile integration test dependencies
* mvn -am -pl tests/integration -DskipTests install
* mvn -am -pl tests/integration -Dcheckstyle.skip=true -Dlicense.skip=true -Dspotbugs.skip=true -DskipTests install
* # compile apachepulsar/java-test-image with async profiler (add "clean" to ensure a clean build with recent changes)
* ./build/build_java_test_image.sh -Ddocker.install.asyncprofiler=true
* ./build/build_java_test_image.sh -Ddocker.install.asyncprofiler=true -Pdocker-wolfi
* # set environment variables
* export PULSAR_TEST_IMAGE_NAME=apachepulsar/java-test-image:latest
* export NETTY_LEAK_DETECTION=off
Expand Down Expand Up @@ -92,31 +98,98 @@ public PulsarPerfContainer(String clusterName,
createContainerCmd.withName(clusterName + "-" + hostname);
});
withEnv("PULSAR_MEM", DEFAULT_PULSAR_MEM);
withEnv("PULSAR_GC", "-XX:+UseZGC -XX:+ZGenerational");
setCommand("sleep 1000000");
File testOutputDir = new File("target");
if (!testOutputDir.exists()) {
if (!testOutputDir.mkdirs()) {
throw new IllegalArgumentException("Test output directory + '" + testOutputDir.getAbsolutePath()
+ "' doesn't exist and cannot be created.");
}
}
if (!testOutputDir.isDirectory()) {
throw new IllegalArgumentException(
"Test output directory '" + testOutputDir.getAbsolutePath() + "' isn't a directory.");
}
// change access to testOutputDir to allow all access so the the container user can write to it
// This matters only on Linux
try {
Files.setPosixFilePermissions(testOutputDir.toPath(), PosixFilePermissions.fromString("rwxrwxrwx"));
} catch (IOException e) {
throw new UncheckedIOException("Cannot change access to test output directory", e);
}
withFileSystemBind(testOutputDir.getAbsolutePath(), "/testoutput", BindMode.READ_WRITE);
}

public CompletableFuture<Long> consume(String topicName) throws Exception {
return DockerUtils.runCommandAsyncWithLogging(getDockerClient(), getContainerId(),
"/pulsar/bin/pulsar-perf", "consume", topicName,
"-u", "pulsar://" + brokerHostname + ":6650",
"-st", "Shared",
"-aq",
"-m", String.valueOf(numberOfMessages), "-ml", "400M");
"bash", "-c", "echo $$ > /tmp/command.pid; "
+ "/pulsar/bin/pulsar-perf consume " + topicName + " "
+ "-u pulsar://" + brokerHostname + ":6650 "
+ "-st Shared "
+ "-q 50000 "
+ "-m " + numberOfMessages + " -ml 400M "
+ "--histogram-file=/testoutput/consume.histogram.$(date +%s).hdr "
+ "2>&1 | tee /testoutput/consume.$(date +%s).txt");
}

public CompletableFuture<Long> produce(String topicName) throws Exception {
return DockerUtils.runCommandAsyncWithLogging(getDockerClient(), getContainerId(),
"/pulsar/bin/pulsar-perf", "produce", topicName,
"-u", "pulsar://" + brokerHostname + ":6650",
"-au", "http://" + brokerHostname + ":8080",
"-r", String.valueOf(Integer.MAX_VALUE), // max-rate
"-s", "8192", // 8kB message size
"-m", String.valueOf(numberOfMessages), "-ml", "400M");
"bash", "-c", "echo $$ > /tmp/command.pid; "
+ "/pulsar/bin/pulsar-perf produce " + topicName + " "
+ "-u pulsar://" + brokerHostname + ":6650 "
+ "-au http://" + brokerHostname + ":8080 "
+ "-r " + Integer.MAX_VALUE + " "
+ "-s 128 -db "
+ "-o 20000 "
+ "-m " + numberOfMessages + " -ml 400M "
+ "--histogram-file=/testoutput/produce.histogram.$(date +%s).hdr "
+ "2>&1 | tee /testoutput/produce.$(date +%s).txt");
}

public CompletableFuture<Long> stats(String topicName) throws Exception {
String basePath = "http://" + brokerHostname + ":8080/admin/v2/" + topicName.replace("://", "/");
// print out stats and internal stats every 10 seconds
return DockerUtils.runCommandAsyncWithLogging(getDockerClient(), getContainerId(),
"bash", "-c",
String.format("echo $$ > /tmp/command.pid; "
+ "while [[ 1 ]]; do "
+ "curl -s %s/stats | jq | tee /testoutput/stats.$(date +%%s).txt; "
+ "sleep 1; "
+ "curl -s %s/internalStats | jq | tee /testoutput/internal_stats.$(date +%%s).txt; "
+ "curl -s http://%s:8080/metrics/ > /testoutput/metrics.$(date +%%s).txt; "
+ " sleep 10; "
+ "done",
basePath, basePath, brokerHostname));
}

public void triggerShutdown() {
if (isRunning()) {
// attempt to stop containers gracefully
DockerUtils.runCommandAsyncWithLogging(getDockerClient(), getContainerId(),
"bash", "-c", "pkill java; while pgrep -c java; do "
+ "echo Waiting for java processes to stop.; sleep 1; done; "
+ "kill $(cat /tmp/command.pid)")
.orTimeout(10, TimeUnit.SECONDS)
.exceptionally(t -> null)
.join();
}
}

public void stop() {
if (isRunning()) {
// attempt to stop containers gracefully
dockerClient.stopContainerCmd(getContainerId())
.withTimeout(15)
.exec();
}
super.stop();
}
}

private PulsarPerfContainer perfConsume;
private PulsarPerfContainer perfProduce;
private PulsarPerfContainer printStats;

@Override
public void setupCluster() throws Exception {
Expand All @@ -126,14 +199,27 @@ public void setupCluster() throws Exception {

@Override
public void tearDownCluster() throws Exception {
if (printStats != null) {
printStats.triggerShutdown();
}
if (perfProduce != null) {
perfProduce.triggerShutdown();
}
if (perfConsume != null) {
perfConsume.stop();
perfConsume = null;
perfConsume.triggerShutdown();
}
if (printStats != null) {
printStats.stop();
printStats = null;
}
if (perfProduce != null) {
perfProduce.stop();
perfProduce = null;
}
if (perfConsume != null) {
perfConsume.stop();
perfConsume = null;
}
super.tearDownCluster();
}

Expand All @@ -142,7 +228,10 @@ protected void beforeStartCluster() throws Exception {
super.beforeStartCluster();
pulsarCluster.forEachContainer(
// This is effective only when -Pdocker-wolfi has been passed when building java-test-image
c -> c.withEnv("GLIBC_TUNABLES", "glibc.malloc.hugetlb=1:glibc.malloc.mmap_threshold=2097152"));
// setting mmap_threshold explicitly will avoid it's dynamic increase
// https://sourceware.org/glibc/manual/latest/html_node/Memory-Allocation-Tunables.html
c -> c.withEnv("GLIBC_TUNABLES",
"glibc.malloc.hugetlb=1:glibc.malloc.mmap_threshold=131072:glibc.malloc.arena_max=4"));
}

@Override
Expand All @@ -160,15 +249,25 @@ protected PulsarClusterSpec.PulsarClusterSpecBuilder beforeSetupCluster(String c
specBuilder.numProxies(0);

// Increase memory for brokers and configure more aggressive rollover
specBuilder.brokerEnvs(Map.of("PULSAR_MEM", BROKER_PULSAR_MEM,
"managedLedgerMinLedgerRolloverTimeMinutes", "1",
"managedLedgerMaxLedgerRolloverTimeMinutes", "5",
"managedLedgerMaxSizePerLedgerMbytes", "512",
"managedLedgerDefaultEnsembleSize", "1",
"managedLedgerDefaultWriteQuorum", "1",
"managedLedgerDefaultAckQuorum", "1",
"maxPendingPublishRequestsPerConnection", "100000"
));
Map<String, String> brokerEnvs = new HashMap<>();
brokerEnvs.put("PULSAR_MEM", BROKER_PULSAR_MEM);
brokerEnvs.put("managedLedgerMinLedgerRolloverTimeMinutes", "1");
brokerEnvs.put("managedLedgerMaxLedgerRolloverTimeMinutes", "5");
brokerEnvs.put("managedLedgerMaxSizePerLedgerMbytes", "512");
brokerEnvs.put("managedLedgerDefaultEnsembleSize", "1");
brokerEnvs.put("managedLedgerDefaultWriteQuorum", "1");
brokerEnvs.put("managedLedgerDefaultAckQuorum", "1");
//brokerEnvs.put("maxPendingPublishRequestsPerConnection", "1000");
brokerEnvs.put("dispatcherRetryBackoffInitialTimeInMs", "0");
brokerEnvs.put("dispatcherRetryBackoffMaxTimeInMs", "0");
brokerEnvs.put("preciseDispatcherFlowControl", "true");
//brokerEnvs.put("PULSAR_PREFIX_subscriptionKeySharedUseClassicPersistentImplementation", "true");
//brokerEnvs.put("PULSAR_PREFIX_subscriptionSharedUseClassicPersistentImplementation", "true");
brokerEnvs.put("dispatcherMaxReadBatchSize", "1000");
//brokerEnvs.put("dispatcherMaxReadSizeBytes", "10000000");
//brokerEnvs.put("dispatcherDispatchMessagesInSubscriptionThread", "false");
//brokerEnvs.put("dispatcherMaxRoundRobinBatchSize", "1000");
specBuilder.brokerEnvs(brokerEnvs);

// Increase memory for bookkeepers and make compaction run more often
Map<String, String> bkEnv = new HashMap<>();
Expand All @@ -190,9 +289,11 @@ protected PulsarClusterSpec.PulsarClusterSpecBuilder beforeSetupCluster(String c
String brokerHostname = clusterName + "-pulsar-broker-0";
perfProduce = new PulsarPerfContainer(clusterName, brokerHostname, "perf-produce");
perfConsume = new PulsarPerfContainer(clusterName, brokerHostname, "perf-consume");
printStats = new PulsarPerfContainer(clusterName, brokerHostname, "print-stats");
specBuilder.externalServices(Map.of(
"pulsar-produce", perfProduce,
"pulsar-consume", perfConsume
"pulsar-consume", perfConsume,
"print-stats", printStats
));

return specBuilder;
Expand All @@ -204,6 +305,8 @@ public void runPulsarPerf() throws Exception {
CompletableFuture<Long> consumeFuture = perfConsume.consume(topicName);
Thread.sleep(1000);
CompletableFuture<Long> produceFuture = perfProduce.produce(topicName);
Thread.sleep(4000);
printStats.stats(topicName);
FutureUtil.waitForAll(List.of(consumeFuture, produceFuture))
.orTimeout(3, TimeUnit.MINUTES)
.exceptionally(t -> {
Expand Down
Loading