【发布时间】:2018-11-06 21:31:20
【问题描述】:
我正在尝试编写一个简单的 JAVA 程序,它会生成一些数据(只是一个 POJO),这些数据会在 Kafka 主题中发布。从这个主题中,订阅者获取数据并将其写入 Cassandra DB。
生产和获取工作正常,但在将数据写入 Cassandra DB 时,有些事情让我感到疑惑。
当我尝试写入数据时,我总是需要打开一个到数据库的新连接。看起来很不愉快。
@Override
public void run() {
setRunning(true);
try {
konsument.subscribe(Collections.singletonList(ServerKonfiguration.TOPIC));
while (running) {
ConsumerRecords<Long, SensorDaten> sensorDaten = konsument.poll(Long.MAX_VALUE);
sensorDaten.forEach(
datum -> {
CassandraConnector cassandraConnector = new CassandraConnector();
cassandraConnector.schreibeSensorDaten(datum.key(), datum.value());
System.out.printf(
"Consumer Record:(%d, %s, %d, %d)\n",
datum.key(), datum.value(), datum.partition(), datum.offset());
});
}
} catch (Exception e) {
e.printStackTrace();
} finally {
konsument.close();
}
}
上面的代码 sn-p 正在工作,但正如我所提到的,对于每次写入,我都必须创建一个新连接。
当我在循环外初始化cassandraConnector 时,我进行了一次成功写入,然后我得到“没有可用的主机”异常。
CassandraConnector 类:
public class CassandraConnector {
private final String KEYSPACE = "ba2";
private final String SERVER_IP = "127.0.0.1";
private Cluster cluster;
private Session session;
public CassandraConnector() {
cluster = Cluster.builder().addContactPoint(SERVER_IP).build();
session = cluster.connect(KEYSPACE);
}
public void schreibeSensorDaten(Long key, SensorDaten datum) {
try {
session.execute(
"INSERT INTO.....
【问题讨论】:
-
为什么不尝试在这种情况下使用 Spark 流,火花流与从 Kafka 读取然后写入 Cassandra 或 Hbase 相关
-
@songoku1610 感谢您的意见。之前没有听说过 Spark。不幸的是,我只被“允许”使用 Kafka 和 Cassandra。
标签: java database cassandra connection