-
Notifications
You must be signed in to change notification settings - Fork 1.9k
/
MessageDAOCassImpl.java
93 lines (81 loc) · 3.12 KB
/
MessageDAOCassImpl.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
package fi.markoa.tfb.servlet3;
import com.datastax.driver.core.*;
import com.google.common.base.Function;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.ListeningExecutorService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.io.InputStream;
import java.util.*;
/**
* Cassandra data access implementation class for the "World" domain model.
*
* @author marko asplund
*/
public class MessageDAOCassImpl implements MessageDAO {
private static final Logger LOGGER = LoggerFactory.getLogger(MessageDAOCassImpl.class);
private static final String CONFIG_FILE_NAME= "/application.properties";
private Cluster cluster;
private Session session;
private Map<String, PreparedStatement> statements;
@Override
public void init(ListeningExecutorService executorService) {
LOGGER.debug("init()");
Properties conf;
try (InputStream is = this.getClass().getClassLoader().getResourceAsStream(CONFIG_FILE_NAME)) {
if(is == null)
throw new IOException("file not found: "+CONFIG_FILE_NAME);
conf = new Properties();
conf.load(is);
} catch (IOException ex) {
LOGGER.error("failed to open config file", ex);
throw new RuntimeException(ex);
}
cluster = Cluster.builder()
.addContactPoint(conf.getProperty("cassandra.host"))
// .withCredentials(conf.getProperty("cassandra.user"), conf.getProperty("cassandra.pwd"))
.build();
session = cluster.connect(conf.getProperty("cassandra.keyspace"));
Map<String, PreparedStatement> stmts = new HashMap<>();
stmts.put("get_by_id", session.prepare("SELECT randomnumber FROM world WHERE id=?"));
stmts.put("update_by_id", session.prepare("UPDATE world SET randomnumber=? WHERE id=?"));
statements = Collections.unmodifiableMap(stmts);
}
@Override
public ListenableFuture<World> read(final int id) {
Function<ResultSet, World> transformation = new Function<ResultSet, World>() {
@Override
public World apply(ResultSet results) {
Row r = results.one();
return new World(id, r.getInt("randomnumber"));
}
};
return Futures.transform(session.executeAsync(statements.get("get_by_id").bind(id)), transformation);
}
public ListenableFuture<List<World>> read(List<Integer> ids) {
List<ListenableFuture<World>> futures = new ArrayList<>();
for(Integer id : ids)
futures.add(read(id));
return Futures.allAsList(futures);
}
public ListenableFuture<Void> update(List<World> worlds) {
Function<ResultSet, Void> transformation = new Function<ResultSet, Void>() {
@Override
public Void apply(ResultSet rows) {
return null;
}
};
BatchStatement bs = new BatchStatement(BatchStatement.Type.UNLOGGED);
for(World w : worlds)
bs.add(statements.get("update_by_id").bind(w.getId(), w.getRandomNumber()));
return Futures.transform(session.executeAsync(bs), transformation);
}
@Override
public void destroy() {
LOGGER.debug("destroy()");
session.close();
cluster.close();
}
}