【发布时间】:2019-02-04 10:19:41
【问题描述】:
我在 HDFS 上有一个镶木地板文件。它每天都会被一个新的覆盖。我的目标是使用 DataStream API 在 Flink Job 中连续发出这个 parquet 文件 - 当它发生变化时 - 作为 DataStream。 最终目标是在广播状态下使用文件内容,但这超出了本问题的范围。
- 要连续处理文件,有一个非常有用的 API:Data-sources about datasources。更具体地说,FileProcessingMode.PROCESS_CONTINUOUSLY:这正是我所需要的。这适用于读取/监控文本文件,没问题,但不适用于镶木地板文件:
// 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)
- 要处理 parquet 文件,我可以通过以下 API 使用 Hadoop 输入格式: using-hadoop-inputformats。但是这个 API 没有 FileProcessingMode 参数,并且只处理一次文件:
// Partial version 2: the parquet file is only processed once
val parquetPath: String = "/path_to_file_dir/parquet_0000"
// raw text format
val hadoopInputFormat: HadoopInputFormat[Void, ArrayWritable] = HadoopInputs.readHadoopFile(new MapredParquetInputFormat(), classOf[Void], classOf[ArrayWritable], parquetPath)
val stream: DataStream[(Void, ArrayWritable)] = streamExecutionEnvironment.createInput(hadoopInputFormat).map { record =>
// process the record here ...
}
我想以某种方式结合这两个 API,通过 DataStream API 连续处理 Parquet 文件。你们有没有人尝试过这样的事情?
【问题讨论】:
标签: scala apache-flink parquet