【发布时间】:2015-11-21 21:19:50
【问题描述】:
我目前正在使用这个堆栈:
- Cassandra 2.2(多节点)
- Spark/Streaming 1.4.1
- Spark-Cassandra-Connector 1.4.0-M3
我有这个 DStream[Ids] ,其中 RDD 大约有 6000-7000 个元素。 id 是分区键。
val ids: DStream[Ids] = ...
ids.joinWithCassandraTable(keyspace, tableName, joinColumns = SomeColumns("id"))
随着tableName 变大,假设大约 30k“行”,查询需要更长的时间,而且我无法保持在批处理持续时间阈值以下。它的执行类似于使用大量的IN-clause,我知道这是不可取的。
有没有更有效的方法来做到这一点?
回答:
在与 Cassandra 进行连接之前,请务必记住使用 repartitionByCassandraReplica 重新分区您的本地 RDD,以确保每个分区仅针对本地 Cassandra 节点工作。就我而言,我还必须在加入本地 RDD/DStream 时增加分区,以使任务在工作人员之间均匀分布。
【问题讨论】:
标签: scala cassandra apache-spark spark-streaming spark-cassandra-connector