CoGroupByKey, emit after n records have been grouped

Viewed 102

I'm trying to parallelize some heavier computations like so:

inputs = (p | "Read" >> beam.io.ReadFromAvro('/mypath/myavrofiles*')
            | "Generate Key" >> beam.Map(lambda row: (gen_key(row), row)))

calc1_results = inputs | "perform calc1" >> beam.Pardo(Calc1())
calc2_results = inputs | "perform calc2" >> beam.Pardo(Calc2())

combined = (({"calc1": calc1_results, "calc2": calc2_results})
            | beam.CoGroupByKey()
            | beam.Values())

final = combined | "Use Grouped results" >> beam.ParDo(PerformFinalCalculation())
  • Each heavy calc emits (key, result)
  • Each key is unique for each input. One input, One result, one Key

Is there some way to emit from the CoGroupByKey after a single result1/result2 has been collected for each key?

Ultimately I'd like to achieve something along the lines of:


               +------------+
               |            |
               |   Input    |
               |            +-----------------+
               +------------+                 |
                     |                        |
          v-------------------v               |
  +------------+       +------------+         |
  |            |       |            |         |
  |  Heavy     |       |   Heavy    |         |
  |  Calc 1    |       |   Calc 2   |         |
  |            |       |            |         |
  +------------+       +------------+         |
             |            |                   |
             |            |                   |
             |            |                   |
          +--v------------v--+                |
          |      Merged      |                |
          |  original dict,  +<---------------+
          |result 1, result2 |
          |                  |
          +------------------+

0 Answers
Related