【问题标题】:What is best approach to join data in spark streaming application?在火花流应用程序中加入数据的最佳方法是什么?
【发布时间】: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


【解决方案1】:

如果您使用的是 Apache Cassandra,那么您只有一种方法可以有效地连接 Cassandra 中的数据 - 通过RDD API's joinWithCassandraTable。 Spark Cassandra 连接器 (SCC) 的开源版本仅支持它,而在 DSE 版本中,有一个代码允许对 Cassandra 执行有效连接,也适用于 Spark SQL - 所谓的DSE Direct Join。如果您在 Spark SQL 中针对 Cassandra 表使用 join,Spark 将需要从 Cassandra 读取所有数据,然后执行连接 - 这非常慢。

我没有 OSS SCC 为 Spark Structured Streaming 进行连接的示例,但我有一些“正常”连接的示例,例如 this

CassandraJavaPairRDD<Tuple1<Integer>, Tuple2<Integer, String>> joinedRDD =
     trdd.joinWithCassandraTable("test", "jtest",
     someColumns("id", "v"), someColumns("id"),
     mapRowToTuple(Integer.class, String.class), mapTupleToRow(Integer.class));

【讨论】:

  • 是 - joinWithCassandraTables 从微批处理中获取记录并使用该记录中的主键从 Cassandra 获取实际记录
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-11-12
  • 1970-01-01
  • 2018-07-08
  • 2019-07-07
  • 2021-11-13
  • 2010-12-30
相关资源
最近更新 更多