1
votes

The Problem

I would like to multiply 2 sparse matrices efficiently under the Spark infrastructure in scalable manner, with the assumption that both matrices can fit into memory.

May Approach

At the beginning, for getting a comparable baseline, I took a dataframe with ~100,000 sparse vectors and performed trivial inner multiplication on one machine (by transforming to scipy csr matrix, spliting to batches and multiple all batches serially). The overall running time I got was 5 minutes.

Next, I implemented a fully map-reduce based algorithm over spark, using CoordinateMatrix as explained here and got a really bad performance - over 50 minutes (!) for the whole multiplication:

def coordinate_matrix_mul(mat_left, mat_right):
    mat_left_cols = mat_left.entries.map(lambda entry: (entry.j, (entry.i, entry.value)))
    mat_right_rows = mat_right.entries.map(lambda entry: (entry.i, (entry.j, entry.value)))

    product_entries = mat_left_cols.join(mat_right_rows)\
                .map(lambda pair: ((pair[1][0][0], pair[1][1][0]), pair[1][0][1]*pair[1][1][1]))\
                .reduceByKey(lambda x,y: x+y)\
                .map(lambda cell: MatrixEntry(cell[0][0], cell[0][1], cell[1])) 

    return CoordinateMatrix(product_entries)

My final approach was to take my first naive attempt (that didn't utilized spark with running time of 5 minutes), and parallelize the batches multiplication by the following steps:

  • broadcast the left and right matrix to the workers.
  • parallel the batches indicies in the form of [(from_row, to_row), ...] of the left matrix to rdd.
  • map over the batches indicies and perform the multiplication simultaneously.

However I got the worst running time between all techniques above - over 60 minutes.

I tried changing the batche size, and checked if there is any lack of memory that causes swaping and threshing, but it wasn't the case.

Anyone have an idea what I am missing?

Thanks in advance.

1
Can you please perform another test with this Spark SQL version: github.com/Fokko/spark-matrix-multiplication - cronoik

1 Answers

0
votes
from pyspark.sql import SparkSession, types, Window, functions as F

a = spark.createDataFrame([
    [0, 0, 1],
    [0, 1, 2],
    [1, 1, 4],
    [2, 0, 5],
    [2, 1, 6]], ['a_row', 'a_column', 'a_value']).cache()

b = spark.createDataFrame([
    [0, 0, 1],
    [0, 1, 2],
    [1, 0, 3],
    [1, 1, 4]], ['b_row', 'b_column', 'b_value']).cache()

c = a.join(b, a.a_column == b.b_row) \
    .withColumn('product', F.udf(lambda x, y: x * y, types.IntegerType())('a_value', 'b_value')) \
    .groupBy(['a_row', 'b_column']).agg(F.sum('product').alias('c_value')) \
    .select(F.column('a_row').alias('c_row'), F.column('b_column').alias('c_column'), 'c_value').cache()

c.show()

+-----+--------+-------+
|c_row|c_column|c_value|
+-----+--------+-------+
|    1|       0|     15|
|    1|       1|     22|
|    0|       1|     10|
|    2|       0|     23|
|    0|       0|      7|
|    2|       1|     34|
+-----+--------+-------+