【问题标题】:Flink forward files from List<String> filePathsFlink 从 List<String> filePaths 转发文件
【发布时间】:2020-08-12 16:16:04
【问题描述】:

我们有一个来自数据库表的文件路径列表,其中包含创建时间的时间戳。试图弄清楚我们如何使用 db 中的文件路径列表仅将那些文件从 nfs 转发到 kafka sink。

现在我正在使用具有文件夹根目录的 ContinuousFileMonitoringFunction 的自定义版本,该文件夹将包含 DB 将显示的所有文件。由于文件夹太大而只有几 TB 的数据,因此此操作非常缓慢,因为要遍历文件夹以收集有关更新文件的信息。

Table orders = tableEnv.from("Customers");
Table result = orders.where($("b").isEqual("****"));

DataSet<String> ds  = result.toDataSet();

ds 包含所有应该发送到 kafka 的文件路径。

以下是我计划实施的想法。但是考虑到 flink 并行性、flink 库支持等,有没有更有效的方法?

public class FileContentMap extends RichFlatMapFunction<String, String> {

      

    @Override
    public void flatMap(String input, Collector<String> out) throws Exception {

       
       
        // get the file path
        String filePath = input;

        String fileContent = readFile(input);

    out.collect(fileCOntent);

       
    }

    @Override
    public void open(Configuration config) {
       
    }
}

DataSet<String> contectDataSet = ds.map(new FileCOntentMap());

contectDataSet.addSink(kafkaProducer);

【问题讨论】:

    标签: apache-flink flink-streaming nfs flink-sql flink-batch


    【解决方案1】:

    你的方法对我来说似乎很好。也许更有效的是创建一个RichParallelSourceFunction,在open() 方法中,您调用数据库以获取已更新的文件列表,并构建一个内存中的文件列表特定的源子任务(类似于filePath.hashCode() % numSubTasks == mySubTask)应该发出以由您的FileContentMap 处理。

    【讨论】:

      猜你喜欢
      • 2015-12-16
      • 1970-01-01
      • 1970-01-01
      • 2021-09-17
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-02-28
      • 1970-01-01
      相关资源
      最近更新 更多