pyspark - preprocessing with a kind of "product-join"

Viewed 31

I have 2 datasets that I can represent as:

  1. The first dataframe is my raw data. It contains millions of row and around 6000 areas.
+--------+------+------+-----+-----+
|  user  | area | time | foo | bar |
+--------+------+------+-----+-----+
| Alice  | A    |    5 | ... | ... |
| Alice  | B    |   12 | ... | ... |
| Bob    | A    |    2 | ... | ... |
| Charly | C    |    8 | ... | ... |
+--------+------+------+-----+-----+
  1. This second dataframe is a mapping table. It has around 200 areas (not 5000) for 150 places. Each area can have 1-N places (and a place can have 1-N areas too). It can be represented unpivoted this way:
+------+--------+-------+
| area | place  | value |
+------+--------+-------+
| A    | placeZ |   0.1 |
| B    | placeB |   0.6 |
| B    | placeC |   0.4 |
| C    | placeA |   0.1 |
| C    | placeB |  0.04 |
| D    | placeA |   0.4 |
| D    | placeC |   0.6 |
| ...  | ...    |   ... |
+------+--------+-------+

or pivoted

+------+--------+--------+--------+-----+
| area | placeA | placeB | placeC | ... |
+------+--------+--------+--------+-----+
| A    |      0 |      0 |      0 | ... |
| B    |      0 |    0.6 |    0.4 | ... |
| C    |    0.1 |   0.04 |      0 | ... |
| D    |    0.4 |      0 |    0.6 | ... |
+------+--------+--------+--------+-----+

I would like to create a kind of product-join to have something like:

+--------+--------+--------+--------+-----+--------+
|  user  | placeA | placeB | placeC | ... | placeZ |
+--------+--------+--------+--------+-----+--------+
| Alice  |      0 |    7.2 |    4.8 |   0 |    0.5 | <- 7.2 & 4.8 comes from area B and 0.5 from area A
| Bob    |      0 |      0 |      0 |   0 |    0.2 |
| Charly |    0.8 |   0.32 |      0 |   0 |      0 |
+--------+--------+--------+--------+-----+--------+

I see 2 options so far:

    • Perform a left join between the main table and the pivoted one
    • Multiply each column by the time (around 150 columns)
    • Groupby user with a sum
    • Perform a outer join between the main table and the unpivoted one
    • Multiply the time by value
    • Pivot place
    • Groupby user with a sum

I don't like the first option because of the number of multiplications involved (the mapping dataframe is quite sparse). I prefer the second option but I see two problems :

  • If someday, the dataset does not have a place represented, the column will not exist and the dataset will have a different shape (hence failing).
  • Some other features like foo, bar will be duplicated with the outer-join and I'll have to handle it on case by case at the grouping stage (sum or average).

I would like to know if there is something more ready-to-use for this kind of product-join in spark ? I have seen the OneHotEncoder but it only provides only a "1" on each column (so it is even worse than the first solution).

Thanks in advance,

Nicolas

0 Answers
Related