FLINK: Read multiple S3 directories one at a time and sort records at end of reading a directory

Viewed 92

I am new to flink, so appreciate any help

Scenario:

I have a bunch of folders in S3. And within each folder I have multiple files that have events corresponding to different groups. The rule is that all the events corresponding to a particular group should be available in a single folder. So I know when I read the files in a folder that I have all the events for the groups within that folder. Now the ask is to sort the events for each group in order.

For example, as per below scenario, when I read the files in Folder1, it contains all the events for group 1 and group 2.

Folders in S3:

Folder 1
|-file1.txt
|-Group1 - Event1
|-Group2 - Event2
|-file2.txt
|-Group1 - Event3
|-Group1 - Event4

Folder 2
|-file3.txt
|-Group3 - Event5

So in order to do this I have the below high level approach in flink

  1. Read the files in a folder and stream them to keyed [by group id] windows
  2. Once done reading all the files, trigger the windows and sort the events through ProcessWindowFunction and sink them to kafka.
  3. Repeat 1 and 2 for the next folder and so on

I would like to know if the above approach is feasible in flink.

  1. How would I know that flink is done reading say folder1 files
  2. After knowing that how can I trigger the windows corresponding to folder1 groups

and also if there is a any other better solution to achieve this

0 Answers
Related