【问题标题】:joinWithCassandraTable getting much slower on growing table sizejoinWithCassandraTable 在增加表大小时变得慢得多
【发布时间】: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


    【解决方案1】:

    “id”是表中的分区键吗?如果没有,我认为它需要这样做,否则您可能正在执行表扫描,随着表变大,运行速度会逐渐变慢。

    此外,为了使用这种方法获得良好的性能,我相信您需要在您的 ids RDD 上使用 repartitionByCassandraReplica() 操作,以便连接是每个节点上的本地操作。

    this

    【讨论】:

    • 确实,现在速度更快了。我正在做repartitionByCassandraReplica(keyspace, tableName),默认有 10 个分区。我只有 2 个执行程序,并且数据正在使用 Murmur3 进行分区。不过,似乎只有一名工作人员正在读取数据,这导致另一个工作人员处于空闲状态。这是分区问题吗?
    • 你只有一个节点吗?通常你会在每个节点上运行一个 spark worker,所以每个 worker 都会加载该节点本地的数据。
    • 不,我有一个驱动程序和两个执行程序,在 3 台机器上。我认为即使默认设置为 10 个分区,它也会尝试分散内容。
    • 如果你有三个节点,那么你需要三个执行器,否则一些数据需要在网络上洗牌,这会减慢速度。如果您正在寻找最快的性能,您希望将所有操作保持在每个节点的本地。
    • 嗯,是的。司机是执行人。我的意思是我有 2 名工人、3 名执行者和 3 台机器。我真的没有为洗牌而苦苦挣扎。无论如何,我将输入 DStream 上的分区提高到repartitionByCassandraReplica,它似乎将工作均匀地分配给了其他工作人员。现在我只需要优化分区,但那是另一个线程。
    猜你喜欢
    • 2018-09-06
    • 1970-01-01
    • 1970-01-01
    • 2017-02-10
    • 2015-01-25
    • 1970-01-01
    • 1970-01-01
    • 2012-03-07
    • 2023-03-12
    相关资源
    最近更新 更多