【问题标题】:Read data from Cassandra for processing in Flink从 Cassandra 读取数据以在 Flink 中处理
【发布时间】:2017-08-21 09:59:41
【问题描述】:

我必须使用 Flink 作为流引擎来处理来自 Kafka 的数据流。为了对数据进行分析,我需要在 Cassandra 中查询一些表。做这个的最好方式是什么?对于这种情况,我一直在寻找 Scala 中的示例。但是我找不到。如何使用 Scala 作为编程语言在 Flink 中读取来自 Cassandra 的数据? Read & write data into cassandra using apache flink Java API 在同一行有另一个问题。它在答案中提到了多种方法。我想知道在我的情况下最好的方法是什么。此外,大多数可用的示例都是用 Java 编写的。我正在寻找 Scala 示例。

【问题讨论】:

    标签: scala cassandra apache-flink


    【解决方案1】:

    我目前在 flink 1.3 中使用 asyncIO 从 cassandra 读取数据。这是它的文档:

    https://ci.apache.org/projects/flink/flink-docs-release-1.3/dev/stream/asyncio.html(如果有 DatabaseClient,您将使用 com.datastax.drive.core.Cluster)

    如果您需要更深入的示例来使用它专门从 cassandra 中读取,请告诉我,但不幸的是我只能提供 java 中的示例。

    编辑 1

    这是我使用 flink 的异步 I/O 从 Cassandra 读取的代码示例。我仍在努力识别和修复一个问题,由于某种原因(无需深入研究),单个查询返回的大量数据,异步数据流的超时被触发,即使它看起来被 Cassandra 很好地返回并且在超时时间之前。但是假设这只是我正在做的其他事情的一个错误,而不是因为这段代码,这对你来说应该可以正常工作(并且对我来说也可以正常工作几个月):

    public class GenericCassandraReader extends RichAsyncFunction<CustomInputObject, ResultSet> {
    
        private final Properties props;
        private Session client;
    
        public GenericCassandraReader(Properties props) {
            super();
            this.props = props;
        }
    
        @Override
        public void open(Configuration parameters) throws Exception {
            client = Cluster.builder()
                    .addContactPoint(props.cassandraUrl)
                    .withPort(props.cassandraPort)
                    .build()
                    .connect(props.cassandraKeyspace);
        }
    
        @Override
        public void close() throws Exception {
            client.close();
        }
    
        @Override
        public void asyncInvoke(final CustomInputObject customInputObject, final AsyncCollector<ResultSet> asyncCollector) throws Exception {
    
            String queryString = "select * from table where fieldToFilterBy='" + customInputObject.id() + "';";
    
            ListenableFuture<ResultSet> resultSetFuture = client.executeAsync(queryString);
    
            Futures.addCallback(resultSetFuture, new FutureCallback<ResultSet>() {
    
                public void onSuccess(ResultSet resultSet) {
                    asyncCollector.collect(Collections.singleton(resultSet));
                }
    
                public void onFailure(Throwable t) {
                    asyncCollector.collect(t);
                }
            });
        }
    }
    

    再次抱歉,耽搁了。希望能解决这个错误,这样我就可以确定了,但在这一点上,有一些参考总比没有好。

    编辑 2

    所以我们最终确定问题不在于代码,而在于网络吞吐量。很多字节试图通过一个不够大的管道来处理它,东西开始备份,一些开始涓涓细流,但是(感谢 datastax cassandra 驱动程序的 QueryLogger 我们可以看到这一点)接收结果所花费的时间每个查询开始爬升到 4 秒,然后是 6 秒,然后是 8 秒,依此类推。

    TL;DR,代码很好,请注意,如果您遇到来自 Flink 的 asyncWaitOperator 的 timeoutExceptions,则可能是网络问题。

    编辑 2.5

    还意识到,由于网络延迟问题,我们最终转而使用 RichMapFunction 来保存我们从 cassandra 读取的数据,这可能是有益的。因此,该作业只需跟踪通过它的所有记录,而不必在每次有新记录通过时都从表中读取以获取其中的所有记录。

    【讨论】:

    • 感谢吉卡尔。我一直在我的 java 代码中使用 Datastax 的客户端(由于需求的变化,我的代码从 Scala 迁移到了 Java)并且它运行良好,尽管我还没有实现 asyncIO。正如您在回答中提到的,您能否提供一个实现 asyncIO 的示例?
    • 实际上我将不得不回复您。我有一个“工作”实例,但最近开始使用大量数据进行测试,现在它抛出了奇怪的 timeoutExceptions(会进入它,但老实说它应该在这里成为它自己的问题)。因此,一旦弄清楚这一点,我将通过更正来编辑答案。
    • @Jicaar 您是否将这种方式视为流(AsyncDataStream)?如果是这样,这段代码多久查询一次 Cassandra?
    • 我有,而且查询非常频繁。我没有关于它查询和获得结果的频率的确切数字,但它跟上了来自 flink 作业的请求量。每秒大约有 1-2 千个请求。 Cassandra 的 metics 表示它收到了一个请求,并在大约 50 毫秒内为 80% 的请求发送了响应(我相信。我已经有一段时间没有查看这些数字了)。
    • @Jicaar 我是 Flink 的新手,我尝试使用 Cassandra 作为源,只是遇到了“Failing the AsyncWaitOperator”异常,你能分享一下 RichMapFunction 是如何完成的吗?顺便说一句,当我只在 Cassandra 中插入 10 行数据时,为什么会发生这个 AsyncWaitOperator 异常?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-11-12
    • 2016-11-08
    • 1970-01-01
    • 1970-01-01
    • 2018-09-12
    • 2016-08-09
    • 2018-02-01
    相关资源
    最近更新 更多