【问题标题】:Flink -- get data from Cassandra as generic ResultSet and convert it to DataSetFlink——从 Cassandra 获取数据作为通用 ResultSet 并将其转换为 DataSet
【发布时间】:2018-09-18 20:27:10
【问题描述】:

我有 StreamExecutionEnvironment 作业,它使用 kafka 简单的 cql 选择查询。 我尝试使用以下代码异步处理此查询:

公共类 GenericCassandraReader 扩展 RichAsyncFunction {

private static final Logger logger = LoggerFactory.getLogger(GenericCassandraReader.class);
private ExecutorService executorService;

private final Properties props;
private Session client;

public ExecutorService getExecutorService() {
    return executorService;
}

public GenericCassandraReader(Properties props, ExecutorService executorService) {
    super();
    this.props = props;
    this.executorService = executorService;
}

@Override
public void open(Configuration parameters) throws Exception {
    client = Cluster.builder().addContactPoint(props.getProperty("cqlHost"))
            .withPort(Integer.parseInt(props.getProperty("cqlPort"))).build()
            .connect(props.getProperty("keyspace"));

}

@Override
public void close() throws Exception {
    client.close();
    synchronized (GenericCassandraReader.class) {
        try {
            if (!getExecutorService().awaitTermination(1000, TimeUnit.MILLISECONDS)) {
                getExecutorService().shutdownNow();
            }
        } catch (InterruptedException e) {
            getExecutorService().shutdownNow();
        }
    }
}

@Override
public void asyncInvoke(final UserDefinedType input, final AsyncCollector<ResultSet> asyncCollector) throws Exception {
    getExecutorService().submit(new Runnable() {
        @Override
        public void run() {
            ListenableFuture<ResultSet> resultSetFuture = client.executeAsync(input.query);

            Futures.addCallback(resultSetFuture, new FutureCallback<ResultSet>() {

                public void onSuccess(ResultSet resultSet) {
                    asyncCollector.collect(Collections.singleton(resultSet));
                }

                public void onFailure(Throwable t) {
                    asyncCollector.collect(t);
                }
            });
        }
    });
}

}

此代码的每个响应都为 Cassandra ResultSet 提供了不同数量的字段。

在 Flink 中处理 Cassandra ResultSet 的任何想法还是我应该使用其他技术来达到我的目标?

提前感谢您的帮助!

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    Cassandra ResultSet 不是线程安全的。最好尝试使用Flink Cassandra connector。或者至少以类似的方式编写您的实现

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2012-09-13
      • 1970-01-01
      • 2023-03-28
      • 1970-01-01
      • 2018-09-12
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多