This plugin
- Configures Jackson for Javaslang serialisation / deserialisation, so Javaslang types can be used as input and output to jax-rs Resources
- Add cyclops-javaslang to your project which enables conversion between JDK / Javaslang / Guava / Jool / simple-react types and adds features such as for-comprehensions to Javaslang.
- Integrates with micro-client, Rest clients (RestClient, AsyncRestClient and NIORestClient) automatically pick up Javaslang mappings
- Adds some reactive programming support for Javaslang Streams namely a. JavaslangReactive mixin allows creating Javaslang reactive-streams Publishers and Subsribers simple b. JavaslangPipes and JavaslangReactive provide integration with simple-react c. Ability to push data into a Javaslang Stream across threads via JavaslangPipes
Simply add this plugin to the classpath - do remember to also add a Server and jax-rs provider (adding micro-grizzly-with-jersey to your classpath too - is a good bet!)
Maven
<dependency>
<groupId>com.aol.microservices</groupId>
<artifactId>micro-javaslang</artifactId>
<version>x.yz</version>
</dependency>
Gradle
compile 'com.aol.microservices:micro-javaslang:x.yz'
Queue queue = QueueFactories.boundedNonBlockingQueue(1000); //set up an Agrona backed wait-free queue
queue.add("world"); //add some data to the queue
JavaslangPipes.register("myQueue",queue); //register the queue
//on a separate thread process data streaming from the queue
JavaslangPipes.stream("myQueue")
.map(this:process)
.forEach(this:save);
@Microserver
@Path("/javaslang")
@Rest
public class JavaslangApp {
public static void main(String[] args){
new MicroserverApp(() -> "javaslang-app").start();
}
@GET
@Produces("application/json")
@Path("/ping")
public ImmutableJavaslangEntity ping() {
return ImmutableJavaslangEntity.builder().value("value")
.list(List.ofAll("hello", "world"))
.mapOfSets(HashMap.<String,Set>empty().put("key1",HashSet.ofAll(Arrays.asList(1, 2, 3))))
.build();
}
@XmlAccessorType(XmlAccessType.FIELD)
@XmlType(name = "")
@XmlRootElement(name = "immutable")
@Getter
@AllArgsConstructor
@Builder
public static class ImmutableJavaslangEntity {
private final String value;
private final List<String> list;
private final Map<String,Set> mapOfSets;
public ImmutableJavaslangEntity() {
this(null,null,null);
}
}
} Running the app and browsing to http://localhost:8080/javaslang-app/javaslang/ping
{
"value": "value",
"list": [
"hello",
"world"
],
"mapOfSets": {
"key1": [
1,
2,
3
]
}
} Javaslang List and JDK8 CompletableFuture
public void cfList(){
CompletableFuture<String> future = CompletableFuture.supplyAsync(this::loadData);
CompletableFuture<List<String>> results1 = Do.add(future)
.add(List.of("first","second"))
.yield( loadedData -> localData-> loadedData + ":" + localData )
.unwrap();
System.out.println(results1.join());
}
private String loadData(){
return "loaded";
}public class ReactiveStreamsPublisherTest implements JavaslangReactive {
@Test
public void publish(){
Stream<Integer> javaslangStream = this.publish(LazyFutureStream.of(1,2,3));
assertThat(javaslangStream.toList(),equalTo(List.ofAll(1,2,3)));
}
} public class ReactiveStreamsSubscriberTest implements JavaslangReactive {
@Test
public void subscribe(){
Stream<Integer> stream = Stream.ofAll(1,2,3);
CyclopsSubscriber<Integer> sub = SequenceM.subscriber();
this.subsribe(stream, sub);
assertThat(sub.sequenceM().toList(),equalTo(Arrays.asList(1,2,3)));
}
}public class ReactiveStreamsAsyncSubscriberTest implements JavaslangReactive {
@Test
public void subscribeAsync(){
Stream<Integer> stream = Stream.ofAll(1,2,3);
CyclopsSubscriber<Integer> sub = SequenceM.subscriber();
this.subsribe(stream, sub,Executors.newFixedThreadPool(1));
assertThat(sub.sequenceM().toList(),equalTo(Arrays.asList(1,2,3)));
}
}