【发布时间】:2021-10-15 12:33:58
【问题描述】:
我有一个并行度为 5 的 flink 工作(现在!!)。并且richFlatMap 流之一在open(Configuration parameters) 方法中打开一个文件。在flatMap操作中没有任何打开动作,它只是读取文件来搜索一些东西。 (有一个实用程序类具有类似utilityClass.searchText('abc') 的方法)。这是样板代码:
public class MyFlatMap extends RichFlatMapFunction<...> {
private MyUtilityFile myFile;
@Override
public void open(Configuration parameters) throws Exception {
myFile.Open("fileLocation");
}
@Override
public void flatMap(...) throws Exception {
String text = myFile.searchText('abc');
if (text != null) // take an action
else // another action
}
}
此文件每天在特定时间由 python 脚本更新。因此,我还应该在 flatMap 流中打开新创建的文件(通过 python 脚本)。
我只是认为这可以由ScheduledExecutorService 完成,只需一个线程池。
我无法在每次 flatMap 调用时都打开这个文件,因为它很大。
这是我正在尝试编写的样板代码:
public class MyFlatMap extends RichFlatMapFunction<...> implements Runnable {
private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
private MyUtilityFile myFile;
@Override
public void run() {
myFile.Open("fileLocation");
}
@Override
public void open(Configuration parameters) throws Exception {
scheduler.scheduleAtFixedRate(this, 1, 1, TimeUnit.HOURS);
myFile.Open("fileLocation");
}
@Override
public void flatMap(...) throws Exception {
String text = myFile.searchText('abc');
if (text != null) // take an action
else // another action
}
}
这个样板文件是否适合 Flink 环境?如果没有,我怎样才能以预定的方式打开文件? (没有“使用kafka更新文件发送事件并通过flink读取事件后”之类的选项)
【问题讨论】:
标签: apache-flink flink-streaming