How to deal with slow Dataflow jobs that uses GroupByKey?

Viewed 313

I have a pretty large table (in TBs) in Big Query and below is the example of the data i have

------------------------------------------------------
|   date       |   id     |   no_items |   unique_id  |
-------------------------------------------------------
|   2018-01-01 |    ifg   |      4     |    AA1       |
|   2018-01-01 |    hui   |      8     |    AA1       |
|   2018-01-01 |    bhi   |      12    |    AA2       |
|   2018-01-02 |    hui   |      3     |    AA1       |
|   2018-01-02 |    ifg   |      10    |    AA1       |
|   2018-01-02 |    bhi   |      5     |    AA2       |
------------------------------------------------------

I have created a beam job that reads and writes into Avro. I am doing GroupByKey using id and unique_id and then comparing the no_items for each of the dates. If the difference between no_items between two dates is more than 5. I am creating a new PCollection with new fields and one extra field that takes the previous of no_items.But the dataflow job is taking long time to complete even with 25 workers and worker_machine_type as 'n1-highcpu-64'

Below is the BEAM pipeline I had written

def sort_grouped_data(element):
    key, value = element
    value = list(value)
    value.sort(key=lambda x: x[2], reverse=True)
    return [element]

def final_segment(element):
    key, value = element
    value = list(value)
    length = len(value)
    for i in range (0, length):
     if abs(value[length - i][1] - value[length - (i + 1)][1]) > 5:
            segment = {"segment_id": segment_id, "segment_start_date": value[length - 1][2],
                       "segment_end_date": value[length - (i + 1)][2], "id": key[1], "unique_id": key[0],
                       "new_value": value[length - i][1]}
    return [segment]


input_or_inputs(
            | "get_key_tuple" >> beam.Map(lambda element: ((element["unique_id"], element["id"]),
                                                           [element["date"], element["no_itmes"]]))
            | "group_by_id_unique_id" >> beam.GroupByKey()
            | "sort_grouped_data" >> beam.ParDo(sort_grouped_data)
            | "final_segment" >> beam.ParDo(final_segment)
)

I feel the dataflow job is slow because of two reasons

  1. Due to the sorting
  2. Due to the comparison between each of the elements, and when the data is huge it is getting slower.

Can anyone help me with this? Is it possible to optimize the code to make it more faster?

0 Answers
Related