【发布时间】:2020-04-16 20:46:52
【问题描述】:
问题:本质上,这意味着,不是为每个流式记录运行 C* 表的连接,而是为火花流中的每个微批处理(微批处理)记录运行连接?
我们几乎最终确定使用 spark-sql 2.4.x 版本,datastax-spark-cassandra-connector 用于 Cassandra-3.x 版本。
但是对于以下场景中的效率有一个基本问题。
对于流数据记录(即 streamingDataSet ),我需要从 Cassandra(C*) 表中查找现有记录(即 cassandraDataset)。
即
Dataset<Row> streamingDataSet = //kafka read dataset
Dataset<Row> cassandraDataset= //loaded from C* table those records loaded earlier from above.
要查找数据,我需要加入上述数据集
即
Dataset<Row> joinDataSet = cassandraDataset.join(cassandraDataset).where(//somelogic)
进一步处理 joinDataSet 以实现业务逻辑 ...
在上述情况下,我的理解是,对于收到的每条记录 从 kafka 流中,它将查询 C* 表,即数据库调用。
如果 C* 表由 数十亿条记录?应该采用什么方法/程序 接下来改进查找 C* 表?
在这种情况下最好的解决方案是什么?我不能从 C* 表并在数据不断添加到 C* 表时进行查找......即 新的查找可能需要新的持久化数据。
如何处理这种情况?任何建议plzz..
【问题讨论】:
-
您使用的是 OSS Cassandra 还是 DataStax Enterprise?
标签: cassandra apache-spark-sql spark-structured-streaming datastax-enterprise spark-cassandra-connector