【发布时间】:2018-06-05 09:58:34
【问题描述】:
我的 flink 程序应该对每个输入记录进行 Cassandra 查找,并根据结果进行一些进一步的处理。
但我目前正忙于从 Cassandra 读取数据。这是我目前想出的代码 sn-p。
ClusterBuilder secureCassandraSinkClusterBuilder = new ClusterBuilder() {
@Override
protected Cluster buildCluster(Cluster.Builder builder) {
return builder.addContactPoints(props.getCassandraClusterUrlAll().split(","))
.withPort(props.getCassandraPort())
.withAuthProvider(new DseGSSAPIAuthProvider("HTTP"))
.withQueryOptions(new QueryOptions().setConsistencyLevel(ConsistencyLevel.LOCAL_QUORUM))
.build();
}
};
for (int i=1; i<5; i++) {
CassandraInputFormat<Tuple2<String, String>> cassandraInputFormat =
new CassandraInputFormat<>("select * from test where id=hello" + i, secureCassandraSinkClusterBuilder);
cassandraInputFormat.configure(null);
cassandraInputFormat.open(null);
Tuple2<String, String> out = new Tuple8<>();
cassandraInputFormat.nextRecord(out);
System.out.println(out);
}
但问题是,每次查找需要将近 10 秒,换句话说,这个 for 循环需要 50 秒才能执行。
如何加快此操作?或者,有没有其他方法可以在 Flink 中查找 Cassandra?
【问题讨论】:
-
您的程序是用于批处理还是流处理?您是作为批处理还是在流中接收输入记录?
-
@avidlearner 程序以流的形式从 Kafka 读取数据。对于我收到的每条记录,我都应该查找 Cassandra。我想出了一个可行的解决方案,我将很快分享它作为答案。但很想知道是否有更有效的方法。
-
然后您可以使用任何 Java 客户端从 Cassandra 获取记录。 Datastax 的客户端可以在处理流时用于 map 或 flatMap 运算符。 CassandraInputFormat 用于将 Cassandra 查询的结果作为 Flink 中的 DataSet 获取。它仅适用于批处理。
-
@avidlearner 你能提供一些例子吗?我搜索了很多,但没有找到:/我现在已经发布了我的答案。
标签: java cassandra apache-flink flink-streaming