@@ -113,7 +113,7 @@ public class ReactiveStreamApiTests {
|
||||
@Test
|
||||
public void continuousRead() {
|
||||
|
||||
Flux<MapRecord<String, String, String>> messages = streamReceiver.receive(fromStart(SensorData.KEY));
|
||||
var messages = streamReceiver.receive(fromStart(SensorData.KEY));
|
||||
|
||||
messages.as(StepVerifier::create)
|
||||
.then(() ->
|
||||
|
||||
@@ -67,10 +67,10 @@ public class SyncStreamApiTests {
|
||||
public void basics() {
|
||||
|
||||
// XADD with fixed id
|
||||
RecordId fixedId1 = streamOps.add(SensorData.RECORD_1234_0);
|
||||
var fixedId1 = streamOps.add(SensorData.RECORD_1234_0);
|
||||
assertThat(fixedId1).isEqualTo(SensorData.RECORD_1234_0.getId());
|
||||
|
||||
RecordId fixedId2 = streamOps.add(SensorData.RECORD_1234_1);
|
||||
var fixedId2 = streamOps.add(SensorData.RECORD_1234_1);
|
||||
assertThat(fixedId2).isEqualTo(SensorData.RECORD_1234_1.getId());
|
||||
|
||||
// XLEN
|
||||
@@ -82,17 +82,17 @@ public class SyncStreamApiTests {
|
||||
}).withMessageContaining("ID specified");
|
||||
|
||||
// XADD with autogenerated id
|
||||
RecordId autogeneratedId = streamOps.add(SensorData.create("1234", "19.8", null));
|
||||
var autogeneratedId = streamOps.add(SensorData.create("1234", "19.8", null));
|
||||
|
||||
assertThat(autogeneratedId.getValue()).endsWith("-0");
|
||||
assertThat(streamOps.size(SensorData.KEY)).isEqualTo(3L);
|
||||
|
||||
// XREAD from start
|
||||
List<MapRecord<String, String, String>> fromStart = streamOps.read(fromStart(SensorData.KEY));
|
||||
var fromStart = streamOps.read(fromStart(SensorData.KEY));
|
||||
assertThat(fromStart).hasSize(3).extracting(MapRecord::getId).containsExactly(fixedId1, fixedId2, autogeneratedId);
|
||||
|
||||
// XREAD resume after
|
||||
List<MapRecord<String, String, String>> fromOffset = streamOps.read(StreamOffset.create(SensorData.KEY, ReadOffset.from(fixedId2)));
|
||||
var fromOffset = streamOps.read(StreamOffset.create(SensorData.KEY, ReadOffset.from(fixedId2)));
|
||||
assertThat(fromOffset).hasSize(1).extracting(MapRecord::getId).containsExactly(autogeneratedId);
|
||||
}
|
||||
|
||||
@@ -104,7 +104,7 @@ public class SyncStreamApiTests {
|
||||
messageListenerContainer.start();
|
||||
}
|
||||
|
||||
CapturingStreamListener streamListener = CapturingStreamListener.create();
|
||||
var streamListener = CapturingStreamListener.create();
|
||||
|
||||
// XREAD BLOCK
|
||||
messageListenerContainer.receive(fromStart(SensorData.KEY), streamListener);
|
||||
|
||||
Reference in New Issue
Block a user