· 7 years ago · Sep 20, 2018, 11:26 AM
1package com.datastax.driver.dse;
2
3import com.datastax.driver.core.Cluster;
4import com.datastax.driver.core.ResultSet;
5import com.datastax.driver.core.ResultSetFuture;
6import com.datastax.driver.core.Row;
7import com.datastax.driver.core.Session;
8import com.google.common.util.concurrent.FutureCallback;
9import com.google.common.util.concurrent.Futures;
10import com.google.common.util.concurrent.ListenableFuture;
11import java.util.UUID;
12import reactor.core.publisher.Flux;
13import reactor.core.publisher.FluxSink;
14import reactor.core.publisher.FluxSink.OverflowStrategy;
15
16public class FluxExamples {
17
18 public static void main(String[] args) {
19 try (Cluster cluster = Cluster.builder().addContactPoint("127.0.1.1").build();
20 Session session = cluster.connect()
21 ) {
22 session.execute(
23 "CREATE KEYSPACE IF NOT EXISTS test WITH replication = { 'class' : 'SimpleStrategy', 'replication_factor' : 1 }");
24 session.execute("CREATE TABLE IF NOT EXISTS test.users (id uuid PRIMARY KEY, name text)");
25 session.execute("TRUNCATE test.users");
26 for (int i = 0; i < 10; i++) {
27 session.execute("INSERT INTO test.users (id, name) VALUES (?, ?)",
28 UUID.randomUUID(),
29 "user" + i);
30 }
31 synchronousExample(session);
32 asynchronousExample(session);
33 }
34 }
35
36 private static void synchronousExample(Session session) {
37 System.out.println("------------ Synchronous Example -----------------");
38 ResultSet rs = session.execute("SELECT * FROM test.users");
39 Long count = Flux.fromIterable(rs)
40 .map(User::new)
41 .doOnNext(System.out::println)
42 .count().block();
43 System.out.printf("Found %d users total%n", count);
44 System.out.println();
45 }
46
47 private static void asynchronousExample(Session session) {
48 System.out.println("------------ Asynchronous Example -----------------");
49 Long count = Flux.<Row>create(sink -> {
50 ResultSetFuture future = session.executeAsync("SELECT * FROM test.users");
51 consumeAndFetchNext(sink, future);
52 }, OverflowStrategy.BUFFER) // ATTENTION can OOM if subscriber not fast enough
53 .map(User::new)
54 .doOnNext(System.out::println)
55 .count().block();
56 System.out.printf("Found %d users total%n", count);
57 System.out.println();
58 }
59
60 private static void consumeAndFetchNext(FluxSink<Row> sink, ListenableFuture<ResultSet> future) {
61
62 Futures.addCallback(future, new FutureCallback<ResultSet>() {
63
64 @Override
65 public void onSuccess(ResultSet rs) {
66 // How far we can go without triggering the blocking fetch:
67 int remainingInPage = rs.getAvailableWithoutFetching();
68 for (int i = 0; i < remainingInPage; i++) {
69 Row row = rs.one();
70 sink.next(row);
71 }
72 boolean wasLastPage = rs.getExecutionInfo().getPagingState() == null;
73 if (wasLastPage) {
74 sink.complete();
75 } else {
76 ListenableFuture<ResultSet> future = rs.fetchMoreResults();
77 consumeAndFetchNext(sink, future);
78 }
79 }
80
81 @Override
82 public void onFailure(Throwable t) {
83 sink.error(t);
84 }
85 });
86 }
87
88
89 public static class User {
90
91 private final UUID id;
92 private final String name;
93
94 public User(Row row) {
95 id = row.getUUID("id");
96 name = row.getString("name");
97 }
98
99 UUID getId() {
100 return id;
101 }
102
103 String getName() {
104 return name;
105 }
106
107 @Override
108 public String toString() {
109 return "User{" +
110 "id=" + id +
111 ", name='" + name + '\'' +
112 '}';
113 }
114 }
115}