【问题标题】:Flink and Cassandra Connection IssueFlink 和 Cassandra 连接问题
【发布时间】:2020-06-02 06:59:07
【问题描述】:

当正常连接在 Flink 的 DataStreams 之外进行时,是否有人遇到过从 Flink 作业连接到 Cassandra 的任何问题?

    Session session = clusterBuilder.getCluster().connect();
    ResultSet resultSet = session.execute(resultStatement.getQuery());

我不是在语言环境中而是在开发环境中面对这个问题。在本地连接中,它工作正常。即使在我将这段代码保存在 DataStream processElement 中时使用相同的 clusterbuilder 设置,连接也会在 Dev 中建立。

我在 main 中遇到 programInvocation 错误,由于 Flink 1.7 的限制,我看不到整个错误。在仪表板中,您无法在 Flink 1.7 中看到整个异常跟踪。作业未提交。

有人对此有任何线索或遇到过类似的事情吗?

【问题讨论】:

  • 您应该可以在 JobManager 的日志中看到整个异常。

标签: java exception cassandra apache-flink data-stream


【解决方案1】:

最可能的原因(我不是 Flink 专家,但我已经看到 Spark 的这个问题)是 Session 对象不可序列化,并且无法发送给 executors/workers。

为了解决这个问题,通常有一个 API 带有显式的 open/close 调用,允许初始化不可序列化的类。如我所见,Flink 有一个 Asynchronous I/O for External Data Access 的概念,它可能用于访问 Cassandra。

【讨论】:

  • 拍得好!我认为这可能是 OP 正在处理的问题。
猜你喜欢
  • 2017-10-19
  • 2020-10-17
  • 1970-01-01
  • 2018-11-28
  • 2018-11-14
  • 1970-01-01
  • 1970-01-01
  • 2017-06-16
相关资源
最近更新 更多