how to get file name of the parquet file during flink data stream

Viewed 144

I have a data stream using parquet input format, I want to get filename of each item. So I can update the file of the record. How can I do it?

DataStream eventStream = streamExecutionEnvironment.readFile(parquetInputFormat, path, FileProcessingMode.PROCESS_CONTINUOUSLY, 20000);

1 Answers

We had to do something like this where we wanted the timestamp that was part of the directory structure, but for batch processing. Our approach was to extend the input format class (HadoopInputFormat in our case), and in the open() call we can use the input split argument to get the file name. Since we were returning a Tuple2<LongWritable, Text>, and the LongWritable (file offset position) wasn't being used, we extracted and stuffed the timestamp into the first field of the result.

I assume you could extend ParquetInputFormat class and do something similar.

Related