Возможен еще вариант написать MapReduce на map'е, сразу в несколько потоков забубенить. Для примера код на C# (но кажется, где-то есть ошибочка):
| Код | using System; using System.Collections.Generic; using System.Text; using System.IO; using System.Threading;
namespace MapReduce { class Program { struct elem { public string input, output; }
const int max_threads = 20; static int count_runs = 0; static Queue<Thread> threads = new Queue<Thread>(); static Queue<elem> paths = new Queue<elem>(); //static Queue<string> input = new Queue<string>(); //static Queue<string> output = new Queue<string>();
static char[] terms = { ' ', '.', ',', '?', '!', '-' }; const string dirTXT = @"z:\TMP\txt\"; //каталог, откуда будут выдираться тексты const string dirDAT = @"z:\TMP\dat\"; //каталог, куда будут складываться обработанные файлы const string dirTMP = @"z:\TMP\tmp\"; //каталог, куда будут временно складываться файлы для обработки
//главный поток обработки текстов static void Main(string[] args) { Console.WriteLine("Starting..."); int i = 0; bool flag = true; while (!Console.KeyAvailable && flag) { FileInfo[] files = (new DirectoryInfo(dirTXT)).GetFiles("*.txt"); int count = 0; foreach (FileInfo f in files) { Thread t = new Thread(start); threads.Enqueue(t); elem a; f.MoveTo(dirTMP + Path.GetFileName(f.FullName)); a.input = f.FullName; a.output = dirDAT + i.ToString() + ".dat"; paths.Enqueue(a); i++; count++; } if (count != 0) Console.WriteLine("Add " + count.ToString() + " files"); //вручную запустим цикл обработки if (threads.Count > 0) { lock (threads) lock (paths) { Thread t = threads.Dequeue(); elem a = paths.Dequeue(); count_runs++; t.Start(a); }
} if (Console.KeyAvailable) flag = Console.ReadKey().Key != ConsoleKey.Escape; } }
/// <summary> /// Обработка файла input, результаты записывем в файл output /// </summary> static void start(object data)//string input, string output) { elem a = (elem)data; StreamReader fin = new StreamReader(a.input); SortedList<string, int> list = new SortedList<string, int>(); //считаем файл в сортированный список while (!fin.EndOfStream) { foreach (string t in fin.ReadLine().Split(terms)) if (list.ContainsKey(t)) list[t]++; else list.Add(t, 1); } //сбросим результат в файл StreamWriter fout = new StreamWriter(a.output); foreach (string t in list.Keys) fout.WriteLine(t + " " + list[t].ToString()); fout.Close(); fin.Close(); FileInfo f = new FileInfo(a.input); f.Delete(); Console.WriteLine(Path.GetFileName(a.input) + " -> " + Path.GetFileName(a.output)); //удалим себя из количества выполняющихся потоков count_runs--; //и скажем, чтобы менеджер попытался запустить далее lock (threads) while ((count_runs < max_threads) && (threads.Count > 0)) { if (threads.Count > 0) { lock (threads) lock (paths) { Thread t = threads.Dequeue(); elem obj = paths.Dequeue(); t.Start(obj); count_runs++; } } } } } }
|
|