【问题标题】:How to read from Cassandra using Apache Flink?如何使用 Apache Flink 从 Cassandra 读取数据?
【发布时间】: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


【解决方案1】:

我想出了一个使用流数据查询 Cassandra 的速度相当快的解决方案。对有同样问题的人有用。

首先,Cassandra 可以用尽可能少的代码进行查询,

Session session = secureCassandraSinkClusterBuilder.getCluster().connect();
ResultSet resultSet = session.execute("SELECT * FROM TABLE");

但问题在于,创建Session 是一项非常耗时的操作,并且每个键空间都应该执行一次。您创建一次Session 并将其重用于所有读取查询。

现在,由于 Session 不是 Java 可序列化的,它不能作为参数传递给像 MapProcessFunction 这样的 Flink 运算符。有几种方法可以解决这个问题,您可以使用 RichFunction 并在其 Open 方法中对其进行初始化,或者使用 Singleton。我将使用第二种解决方案。

如下创建一个单例类,我们在其中创建Session

public class CassandraSessionSingleton {
    private static CassandraSessionSingleton cassandraSessionSingleton = null;

    public Session session;

    private CassandraSessionSingleton(ClusterBuilder clusterBuilder) {
        Cluster cluster = clusterBuilder.getCluster();
        session = cluster.connect();
    }

    public static CassandraSessionSingleton getInstance(ClusterBuilder clusterBuilder) {
        if (cassandraSessionSingleton == null)
            cassandraSessionSingleton = new CassandraSessionSingleton(clusterBuilder);
        return cassandraSessionSingleton;
    }

}

然后,您可以利用此会话进行所有未来的查询。这里我以ProcessFunction 为例进行查询。

public class SomeProcessFunction implements ProcessFunction <Object, ResultSet> {
    ClusterBuilder secureCassandraSinkClusterBuilder;

    // Constructor
    public SomeProcessFunction (ClusterBuilder secureCassandraSinkClusterBuilder) {
        this.secureCassandraSinkClusterBuilder = secureCassandraSinkClusterBuilder;
    }

    @Override
    public void  ProcessElement (Object obj) throws Exception {
        ResultSet resultSet = CassandraLookUp.cassandraLookUp("SELECT * FROM TEST", secureCassandraSinkClusterBuilder);
        return resultSet;
    }
}

请注意,您可以将ClusterBuilder 传递给ProcessFunction,因为它是可序列化的。现在是我们执行查询的cassandraLookUp 方法。

public class CassandraLookUp {
    public static ResultSet cassandraLookUp(String query, ClusterBuilder clusterBuilder) {
        CassandraSessionSingleton cassandraSessionSingleton = CassandraSessionSingleton.getInstance(clusterBuilder);
        Session session = cassandraSessionSingleton.session;
        ResultSet resultSet = session.execute(query);
        return resultSet;
    }
}

单例对象仅在第一次运行查询时创建,之后重复使用相同的对象,因此查找时没有延迟。

【讨论】:

    猜你喜欢
    • 2018-09-12
    • 2017-07-25
    • 2018-09-13
    • 2017-08-21
    • 1970-01-01
    • 1970-01-01
    • 2014-03-26
    • 1970-01-01
    • 2018-05-09
    相关资源
    最近更新 更多