【问题标题】:how to get file name of the parquet file during flink data stream如何在flink数据流期间获取parquet文件的文件名
【发布时间】:2021-08-11 07:49:40
【问题描述】:
我有一个使用 parquet 输入格式的数据流,我想获取每个项目的文件名。所以我可以更新记录的文件。
我该怎么做?
DataStream eventStream = streamExecutionEnvironment.readFile(parquetInputFormat, path, FileProcessingMode.PROCESS_CONTINUOUSLY, 20000);
【问题讨论】:
标签:
apache-flink
parquet
flink-streaming
【解决方案1】:
我们必须做这样的事情,我们想要作为目录结构一部分的时间戳,但用于批处理。我们的方法是扩展输入格式类(在我们的例子中为HadoopInputFormat),并且在open() 调用中,我们可以使用输入拆分参数来获取文件名。由于我们返回的是Tuple2<LongWritable, Text>,而LongWritable(文件偏移位置)没有被使用,我们提取时间戳并将其填充到结果的第一个字段中。
我假设你可以扩展 ParquetInputFormat 类并做类似的事情。