我假设您追求的是最快的方式。因为每个人都知道最快就是最好的;)
一个简单的方法是使用一堆线程来为你写作。
但是,除非您的文件系统可以很好地扩展,否则您不会通过这样做获得太多好处。 (我在基于 Lustre 的集群系统上使用这种技术,在“大量文件”可能意味着 10k 的情况下 - 在这种情况下,许多写入将发送到不同的服务器/磁盘)
代码看起来像这样:(注意我认为这个版本不正确,因为少量文件会填满工作队列 - 但无论如何,请查看下一个版本以获得更好的版本......)
public void call(Iterator<Tuple2<Text, BytesWritable>> arg0) throws Exception {
int nThreads=5;
ExecutorService threadPool = Executors.newFixedThreadPool(nThreads);
ExecutorCompletionService<Void> ecs = new ExecutorCompletionService<>(threadPool);
int nJobs = 0;
while (arg0.hasNext()) {
++nJobs;
final Tuple2<Text, BytesWritable> tuple2 = arg0.next();
ecs.submit(new Callable<Void>() {
@Override Void call() {
System.out.println(tuple2._1().toString());
String path = "/home/suv/junk/sparkOutPut/"+tuple2._1().toString();
try(PrintWriter writer = new PrintWriter(path, "UTF-8") ) {
writer.println(new String(tuple2._2().getBytes()))
}
return null;
}
});
}
for(int i=0; i<nJobs; ++i) {
ecs.take().get();
}
}
更好的办法是在获得第一个文件的数据后立即开始写入文件,而不是在获得所有文件的数据时开始写入文件 - 并且此写入不会阻塞计算线程。
为此,您将应用程序拆分为多个部分,通过一个(线程安全)队列进行通信。
代码最终看起来更像这样:
public void main() {
SomeMultithreadedQueue<Data> queue = ...;
int nGeneratorThreads=1;
int nWriterThreads=5;
int nThreads = nGeneratorThreads + nWriterThreads;
ExecutorService threadPool = Executors.newFixedThreadPool(nThreads);
ExecutorCompletionService<Void> ecs = new ExecutorCompletionService<>(threadPool);
AtomicInteger completedGenerators = new AtomicInteger(0);
// Start some generator threads.
for(int i=0; ++i; i<nGeneratorThreads) {
ecs.submit( () -> {
while(...) {
Data d = ... ;
queue.push(d);
}
if(completedGenerators.incrementAndGet()==nGeneratorThreads) {
queue.push(null);
}
return null;
});
}
// Start some writer threads
for(int i=0; i<nWriterThreads; ++i) {
ecs.submit( () -> {
Data d
while((d = queue.take())!=null) {
String path = data.path();
try(PrintWriter writer = new PrintWriter(path, "UTF-8") ) {
writer.println(new String(data.getBytes()));
}
return null;
}
});
}
for(int i=0; i<nThreads; ++i) {
ecs.take().get();
}
}
请注意,我没有提供队列类的实现,您可以轻松地包装标准的 java 线程安全类来获得所需的内容。
还有很多事情可以做来减少延迟等 - 这是我用来减少时间的一些进一步的事情......
甚至不要等待为给定文件生成所有数据。传递另一个包含要写入的字节包的队列。
注意分配 - 您可以重复使用一些缓冲区。
nio 的内容存在一些延迟 - 您可以通过使用 C 写入和 JNI 以及直接缓冲区来获得一些性能改进。
线程切换可能会受到伤害,队列中的延迟可能会受到伤害,因此您可能需要稍微批量处理数据。用 1 来平衡这一点可能很棘手。