Reading continuous parquet files from Azure Blob Storage as DataStream

Viewed 71
  1. I want to read parquet files that are generated on Azure blob storage every second in the directory having format YYYY-MM-DD/HH/MM/file_name.parquet

  2. I want to read the data as DataStream.

Do we have any Flink source already defined that can do this?

I tried to use :

// Partial version 1: the raw file is processed continuously
val path: String = "hdfs://hostname/path_to_file_dir/"
val textInputFormat: TextInputFormat = new TextInputFormat(new Path(path))
// monitor the file continuously every minute
val stream: DataStream[String] = streamExecutionEnvironment.readFile(textInputFormat, path, FileProcessingMode.PROCESS_CONTINUOUSLY, 60000)

Issues in the above approach:

  1. Was not able to find readParquetFile support.

  2. It was processing file only from a specific directory passed in the path. In my case directories are changing every minute given the directory format is YYYY-MM-DD/HH/MM/file_name.parquet

0 Answers
Related