It's possible to get rid of Cartesian product, without using a special data structure.
I will demonstrate the method by an example with pyspark
Suppose you have a data set of (user, query).
input:
users_queries_rdd = sc.parallelize([
('u1', 'q1'), ('u1', 'q2'), ('u1', 'q3'), ('u1', 'q4'),
('u2', 'q2'), ('u2', 'q4'), ('u2', 'q5'),
('u3', 'q1'), ('u3', 'q2'), ('u3', 'q4')
])
You would like to count the occurrences of 2 queries for different users.
expected output:
[(('q4', 'q2'), 3),
(('q5', 'q2'), 1),
(('q3', 'q1'), 1),
(('q5', 'q4'), 1),
(('q4', 'q1'), 2),
(('q3', 'q2'), 1),
(('q2', 'q1'), 2),
(('q4', 'q3'), 1)]
方法 1 - 使用笛卡尔积:
pair_queries_count_rdd = users_queries_rdd\
.cartesian(users_queries_rdd)\
.filter(lambda line: line[0] > line[1])\
.filter(lambda line: line[0][0] == line[1][0])\
.map(lambda line: (line[0][1], line[1][1]))\
.map(lambda line: (line, 1))\
.reduceByKey(add)
方法 2 - 摆脱笛卡尔积:
pair_queries_count_rdd_no_cartesian = users_queries_rdd\
.map(lambda line: (line[0], [line[1]]))\
.reduceByKey(add)\
.map(lambda line: tuple(combinations(line[1], 2)))\
.flatMap(lambda line: [(x, 1) for x in line])\
.reduceByKey(add)
解释:
方法一:
1.1 .cartesian(users_queries_rdd)
在 rdd 与自身之间创建笛卡尔积。
它会生成 n^2 个组合。
1.2 .filter(lambda line: line[0] > line[1])
如果包含 (q1,q2),则 (q2,q1) 将被过滤掉。
1.3 .filter(lambda line: line[0][0] == line[1][0])
按用户分组查询(聚合)。
1.4 .map(lambda line: (line[0][1], line[1][1]))
省略用户列。只保留查询对。
1.5 .map(lambda line: (line, 1))
每行将 (q[i], q[j]) 映射到 ((q[i], q[j]), 1)
1.6 .reduceByKey(add)
计算每个查询对的出现次数。
方法二:
2.1 .map(lambda line: (line[0], [line[1]]))
将每一行的 (u[i], q[j]) 映射到 (u[i], [q[j]]) (查询被封装到一个列表中)
2.2 .reduceByKey(add)
为每个用户创建一个包含所有查询的列表。
2.3 .map(lambda line: tuple(combinations(line[1], 2)))
对于每个用户,省略用户列,并创建他们所有查询的组合
2.4 .flatMap(lambda line: [(x, 1) for x in line])
平面映射将平面所有键: (q[i],q[j]) 映射到 ((q[i],q[j]), 1)
2.5 .reduceByKey(add)
计算每个查询对的出现次数。
我们假设每个用户的查询数量相对较少(一个常数)。
所以方法2的效率是O(n)
这个假设是必不可少的,因为每个用户都使用了 iter.combinations 函数。
在大多数实际情况下,方法 2 更有效。