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
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:
Was not able to find readParquetFile support.
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