apache beam con python pardo con input multi pCollections

Sep 21 2020

Sono nuovo in apache_beam e sto cercando di sviluppare una pipeline. Ho 2 pCollection con lo stesso formato e un'altra pCollection con un altro formato. Provo a eseguire una funzione ParDo che per ogni elemento in pCollection 3 dipende da un valore o questo elemento cerca se l'elemento esiste in pCollection 1 o 2 per completare l'output con le informazioni di pCollection 1 o 2. Ma non so come fare questa funzione ParDo .

Questo è il mio codice:

output = (
      pCollection1, pCollection2, pCollection3
      | 'ParDo function' >> beam.ParDo(SearchData()))

E questa è la mia funzione ParDo:

class SampleScores(beam.DoFn):
    def process(self,element):
      
      # here I don't know how call a collection because I have only a "element"

      return xxx

Grazie

Risposte

Jorge Sep 23 2020 at 13:35

Risolto.

Hai esaminato gli ingressi laterali? beam.apache.org/documentation/programming-guide/#side- input Se capisco correttamente la tua domanda, quello che vuoi è avere un processo (self, element, pcoll1, pcoll2), gli input secondari potrebbero aiutarti in questo. - Milan Cermak ieri

Grazie @MilanCermak