I've been researching into the reactive programming and cassandra and am relatively new to both technologies. The issue I'm experiencing is that I'm unable to achieve the same insert performance as the cassandra-stress test. When I run the stress test from my machine to the cluster which is in the same network, I get the following result:
Results:
Op rate : 51,391 op/s [WRITE: 51,391 op/s]
Partition rate : 51,391 pk/s [WRITE: 51,391 pk/s]
Row rate : 51,391 row/s [WRITE: 51,391 row/s]
Latency mean : 3.9 ms [WRITE: 3.9 ms]
Latency median : 2.0 ms [WRITE: 2.0 ms]
Latency 95th percentile : 3.1 ms [WRITE: 3.1 ms]
Latency 99th percentile : 20.6 ms [WRITE: 20.6 ms]
Latency 99.9th percentile : 429.4 ms [WRITE: 429.4 ms]
Latency max : 581.4 ms [WRITE: 581.4 ms]
Total partitions : 15,000,000 [WRITE: 15,000,000]
Total errors : 0 [WRITE: 0]
Total GC count : 0
Total GC memory : 0.000 KiB
Total GC time : 0.0 seconds
Avg GC time : NaN ms
StdDev GC time : 0.0 ms
Total operation time : 00:04:51
This proves I should be able to do anywhere between 40k and 50k inserts per second without an issue from my machine. In reality, though, I'm nowhere near that number as I am only able to achieve just about 10k-15k inserts per second.
04:17:03.782 [p-db-i-8] -> 10,000 in 1 seconds | 14,820,000 of 15,000,000
04:17:04.279 [p-db-i-2] -> 10,000 in 0.497 seconds | 14,830,000 of 15,000,000
04:17:04.711 [p-db-i-2] -> 10,000 in 0.432 seconds | 14,840,000 of 15,000,000
04:17:05.143 [p-db-i-4] -> 10,000 in 0.432 seconds | 14,850,000 of 15,000,000
04:17:05.971 [p-db-i-2] -> 10,000 in 0.828 seconds | 14,860,000 of 15,000,000
04:17:06.807 [p-db-i-8] -> 10,000 in 0.835 seconds | 14,870,000 of 15,000,000
04:17:07.619 [p-db-i-2] -> 10,000 in 0.812 seconds | 14,880,000 of 15,000,000
04:17:08.464 [p-db-i-4] -> 10,000 in 0.845 seconds | 14,890,000 of 15,000,000
04:17:08.884 [p-db-i-3] -> 10,000 in 0.420 seconds | 14,900,000 of 15,000,000
04:17:09.762 [p-db-i-3] -> 10,000 in 0.878 seconds | 14,910,000 of 15,000,000
04:17:10.632 [p-db-i-8] -> 10,000 in 0.870 seconds | 14,920,000 of 15,000,000
04:17:11.146 [p-db-i-1] -> 10,000 in 0.514 seconds | 14,930,000 of 15,000,000
04:17:12.139 [p-db-i-3] -> 10,000 in 0.993 seconds | 14,940,000 of 15,000,000
04:17:13.068 [p-db-i-8] -> 10,000 in 0.929 seconds | 14,950,000 of 15,000,000
04:17:13.935 [p-db-i-8] -> 10,000 in 0.867 seconds | 14,960,000 of 15,000,000
04:17:14.364 [p-db-i-6] -> 10,000 in 0.429 seconds | 14,970,000 of 15,000,000
04:17:14.844 [p-db-i-4] -> 10,000 in 0.480 seconds | 14,980,000 of 15,000,000
04:17:15.667 [p-db-i-1] -> 10,000 in 0.823 seconds | 14,990,000 of 15,000,000
04:17:16.599 [p-db-i-8] -> 10,000 in 0.932 seconds | 15,000,000 of 15,000,000
04:17:16.600 [p-db-i-8] Inserted 15,000,000 records in 18.28 minutes (expected: 15,000,000)
And cassandra metrics also records about 10-15k inserts/s:
I'm a bit puzzled as to what's going wrong here. Essentially the idea is to achieve similar insert rate using reactive, but it looks like I'm doing something wrong.
My project is using Spring Boot 2.4.4 and I am only using the Spring Data Cassandra Reactive dependency:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-cassandra-reactive</artifactId>
</dependency>
There aren't any custom configuration classes nor driver configuration customisations. I have enabled the @EnableReactiveCassandraRepositories and am using flux / monos.
Here is the code:
application.properties:
spring.data.cassandra.schema-action=CREATE_IF_NOT_EXISTS
spring.data.cassandra.request.timeout=10s
spring.data.cassandra.request.throttler.type=none
spring.data.cassandra.request.page-size=5000
spring.data.cassandra.connection.connect-timeout=10s
spring.data.cassandra.connection.init-query-timeout=10s
spring.data.cassandra.repositories.type=reactive
spring.data.cassandra.pool.idle-timeout=120
spring.data.cassandra.local-datacenter=dc1
spring.data.cassandra.contactpoints=10.42.0.2:9042,10.42.0.3:9042,10.42.0.4:9042
spring.data.cassandra.request.consistency=local_one
spring.data.cassandra.keyspace-name=testing
spring.data.cassandra.username=cassandra
spring.data.cassandra.password=cassandra
TestEntity:
import org.springframework.data.cassandra.core.cql.PrimaryKeyType;
import org.springframework.data.cassandra.core.mapping.PrimaryKeyColumn;
import org.springframework.data.cassandra.core.mapping.Table;
@Table(TestEntity.TABLE_NAME)
public class TestEntity
{
public final static String TABLE_NAME = "test_table";
@PrimaryKeyColumn(name = "uuid", ordinal = 0, type = PrimaryKeyType.PARTITIONED)
private UUID uuid;
@PrimaryKeyColumn(name = "version", ordinal = 1, type = PrimaryKeyType.PARTITIONED)
private String version;
public TestEntity()
{
this.uuid = UUID.randomUUID();
}
public TestEntity(String version)
{
this.uuid = UUID.randomUUID();
this.version = version;
}
// standard getter/setters, removed for better readability.
}
TestEntityRepository:
import java.util.UUID;
import org.springframework.data.cassandra.repository.Query;
import org.springframework.data.cassandra.repository.ReactiveCassandraRepository;
import org.springframework.stereotype.Repository;
import reactor.core.publisher.Mono;
@Repository
public interface TestEntityRepository
extends ReactiveCassandraRepository<TestEntity, UUID>
{
@Query("TRUNCATE " + TestEntity.TABLE_NAME)
Mono<Void> truncate();
}
ParallelInsertTest:
import java.text.DecimalFormat;
import java.time.Duration;
import java.time.Instant;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import org.junit.jupiter.api.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Scheduler;
import reactor.core.scheduler.Schedulers;
import reactor.test.StepVerifier;
@SpringBootTest
public class ParallelInsertTest
{
private final Logger logger = LoggerFactory.getLogger(ParallelInsertTest.class);
@Autowired
private TestEntityRepository repository;
@Test
public void testInsert()
throws Exception
{
int cpus = Runtime.getRuntime().availableProcessors();
int rate = 1024; // incrementing this value doesn't actually help that much.
int size = 15000000;
logger.info("Initializing scheduler.");
Scheduler parallelScheduler = Schedulers.newParallel("p-db-i", cpus);
AtomicInteger inserted = new AtomicInteger(0);
AtomicReference<Instant> startTime = new AtomicReference<>();
AtomicReference<Instant> checkTime = new AtomicReference<>();
logger.info("Truncating before proceeding..");
repository.truncate().block();
logger.info("Generating stream of {} integers", size);
Flux<TestEntity> list = Flux.range(1, size)
.flatMap((i -> Mono.just(new TestEntity(i.toString()))));
logger.info("Inserting with parallelism {}", parallelScheduler);
Flux<TestEntity> flux;
flux = Flux.from(list)
.parallel(cpus).runOn(parallelScheduler).sequential()
.limitRate(rate, rate / 4)
.flatMap(value -> repository.save(value).publishOn(parallelScheduler))
.doOnNext((v) -> {
inserted.incrementAndGet();
int a = 10000;
if (inserted.get() % a == 0)
{
logger.info("-> {} in {} | insert {} of {}", formatNumber(a), elapsedTime(checkTime.get()), formatNumber(inserted.get()), formatNumber(size));
checkTime.set(Instant.now());
}
})
.doOnSubscribe((s) -> {
Instant now = Instant.now();
startTime.set(now);
checkTime.set(now);
})
.doFinally((s) -> {
logger.info("Inserted {} records in {} (expected: {})", formatNumber(inserted.get()), elapsedTime(startTime.get()), formatNumber(size));
})
.subscribeOn(parallelScheduler)
.publishOn(parallelScheduler)
;
StepVerifier.create(flux)
.thenConsumeWhile((v) -> true)
.thenAwait(Duration.ofSeconds(600))
.expectComplete()
.verifyThenAssertThat()
.hasNotDroppedElements()
.hasNotDroppedErrors()
.hasNotDiscardedElements()
;
}
// just debugging output helpers.
private String elapsedTime(Instant start)
{
Duration duration = Duration.between(start, Instant.now());
String elapsedTime = "" + duration.getSeconds() + " seconds";
if (duration.getSeconds() < 1)
{
elapsedTime = "0." + duration.toMillis() + " seconds";
}
else if (duration.getSeconds() > 60)
{
double minutes = (double) duration.getSeconds() / 60;
elapsedTime = new DecimalFormat("#.00").format(minutes) + " minutes";
}
return elapsedTime;
}
private String formatNumber(int number)
{
return new DecimalFormat("###,###.###").format(number);
}
}
JVM settings:
-Xms4G
-Xmx8G

