【问题标题】:Why is Cassandra so slow with a composite partition key?为什么 Cassandra 使用复合分区键会这么慢?
【发布时间】:2016-11-17 22:10:48
【问题描述】:

这是我正在处理的问题:

我有这个 33 GB 的 tsv 文件,它有 2 列,第一列是 a_id,第二列是逗号分隔的 b_id 集。问题是,我需要能够检索 b_id 的所有 a_id,所以我将文件加载到 Spark 中,我解析它,我将它平面映射并将其插入到由 b_id 分区的 Cassandra 表中。这个过程大约需要 4 个小时,每个分区 10~15 分钟,并加载所有 200 M a_id,平均每个 20 b_id,所以总共大约 4 B 行。

问题是,因为一些 b_id 很常见,其中一些分区非常大,最大的有 170 万个单元格。所以我尝试计算 a_id 上的哈希并向我正在使用的表添加一个新列(我实际上创建了一个新的单独表),转换为复合分区键。结果是写入每个分区所需的时间增加了 6 倍!!

起初,我认为问题在于我在 Spark 中通过内置的 python hash() 进行的哈希计算,所以我用一个更简单的函数替换了它,它只取模最后 20 位a_id 按我想要的“子分区”数量(5),但没有任何改变......

无论如何,我都不是 Cassandra 方面的专家,但对我来说这没有任何意义。为什么会这样?

【问题讨论】:

  • 什么?! Python hash 用于存储到数据库中,伙计,这真的错了hash 保证仅在一个解释器执行期间保持不变。
  • 好的,注意到了,我得改一下
  • 您确定您使用的是复合分区键而不是集群键吗?
  • 是的,我的创建表语句看起来像“CREATE TABLE hashed_table (a_id text, b_id text, hash int, PRIMARY KEY ((b_id, hash), a_id))”
  • 你能给出一个示例行,如 a_id 和 b_id 的样子吗?

标签: python apache-spark cassandra pyspark


【解决方案1】:

如果没有看到您的 PySpark 代码,我不能 100% 确定,但我怀疑减速是因为您使用无法“下推”并在 Spark Worker 的 JVM 中完成的 Python 函数来操作数据。

当您只是在做一个简单的平面地图时(我假设在一个 RDD 上使用 Spark APIs),Spark 能够在 JVM 内完成该功能。但是一旦你开始在这些 API 之外的 Python 中做“自定义”的东西,Spark 必须在 Spark 工作 JVM 和 Python 之间序列化和流式传输你的数据,以便它可以运行你的 Python 代码来操作数据。我相信它会通过一个很慢的套接字来做到这一点。您可以在此处查看有关 PySpark 内部结构的更多信息:

https://cwiki.apache.org/confluence/display/SPARK/PySpark+Internals

【讨论】:

    猜你喜欢
    • 2017-06-11
    • 1970-01-01
    • 2018-11-16
    • 1970-01-01
    • 1970-01-01
    • 2016-10-08
    • 2014-03-01
    • 1970-01-01
    • 2014-09-23
    相关资源
    最近更新 更多