We're using Airflow for job scheduling, and calling Apache Beam for the ETL step. The data source is unstructured files (batch) which need to be parsed before they can be turned into PCollections. It appears to me that the two best options available are:
- Add a preprocessing node to the Airflow DAG to parse the files and write to a parquet file, which is then processed by Beam.
- Write a custom IO connector in Beam to parse the unstructured file and create the PCollection.
Which option better fits Beam best practices?