Output file to dynamic path based on current date - Apache Beam Python SDK

Viewed 47

I am using fileio.WriteToFiles to output stream data to google cloud but the dynamic path is not working.

Current behavior: The current behavior is the output path is always the date I ran the data pipeline on Airflow. It is not calculated dynamically by the current date.

Expected behavior: I want to output files of my stream apache beam pipeline to google cloud based on the current date when the data is processed. For example, today is July 18. The file needs to be saved to gs://folder/7/18/file_name.txt.

The code is like this:

(input
| "Window into fixed intervals" >> beam.WindowInto(FixedWindows(self.window_size))
| "Write to GCS Bucket" >> fileio.WriteToFiles(
                "/".join(['gs://folder', 
                           DateFormat.get_now_for_timezone(date_format='%Y/%m/%d')]),
                shards=1,
                max_writers_per_bundle=0,
                                              )
)
1 Answers

Is the job templated? The filepath of your WriteToFiles transform could be static the first time you created the template, meaning the whole

"/".join(['gs://folder', 
           DateFormat.get_now_for_timezone(date_format='%Y/%m/%d')])

was only executed once during the pipeline construction time.

Instead, you can provide a callable to the parameter file_naming for WriteToFiles.

file_naming (callable): A callable that takes in a window, pane, shard_index, total_shards and compression; and returns a file name.

E.g., you can alter the default_file_naming callable:

def my_default_file_naming(prefix, suffix=None):
  def _inner(window, pane, shard_index, total_shards, compression, destination):
    # Add your date in the callable
    prefix += DateFormat.get_now_for_timezone(date_format='%Y/%m/%d')
    return _format_shard(
        window, pane, shard_index, total_shards, compression, prefix, suffix)

  return _inner
Related