【发布时间】: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