Inner join with streaming data on Apache Beam/Dataflow

Viewed 748

I know it is possible to perform inner joins on data on Apache Beam, but I could not find an answer for my specific problem (it is a little more complicated than just joining data).

I am receiving data from two PCollections, which I need to merge. For example:

This is what I receive from PCollection A:

{ "key": "key1", "score": ... }
{ "key": "key2", "score": ... }
{ "key": "key3", "score": ... }
...

And this is what I receive from PCollection B:

{ "id": ..., "key": "key1", "otherScore": ... }
{ "id": ..., "key": "key1", "otherScore": ... }
{ "id": ..., "key": "key2", "otherScore": ... }
{ "id": ..., "key": "key3", "otherScore": ... }
{ "id": ..., "key": "key3", "otherScore": ... }
{ "id": ..., "key": "key3", "otherScore": ... }
...

I would like to join both streams by "key". My output to be this:

{ "id": ..., "key": "key1", "score": ..., "otherScore": ... }
{ "id": ..., "key": "key1", "score": ..., "otherScore": ... }
{ "id": ..., "key": "key2", "score": ..., "otherScore": ... }
{ "id": ..., "key": "key3", "score": ..., "otherScore": ... }
{ "id": ..., "key": "key3", "score": ..., "otherScore": ... }
{ "id": ..., "key": "key3", "score": ..., "otherScore": ... }
...

This is a bit different to your 1-1 join because:

  • The entries from PCollection A can be used for several PCollection B. I deally, I would like to store those values for about 4 hours and then drop them. This "caching" is very important.

  • The entries from PCollection B can only be used once. They can stay around for about 10 minutes, waiting to be merged, and then can be dropped.

Is this possible to do on Apache Beam? I would appreciate if someone could point me in the right direction.

0 Answers
Related