X-Git-Url: https://git.donarmstrong.com/?p=mothur.git;a=blobdiff_plain;f=countseqscommand.cpp;h=411a814d73671ebb69b43a783e36e31d6263fe45;hp=dfa012eeff0f0f49c778c00b2098348a489f1855;hb=d1c97b8c04bb75faca1e76ffad60b37a4d789d3d;hpb=4b54ce99af7db8019ea907cd7c2edf789369ada9 diff --git a/countseqscommand.cpp b/countseqscommand.cpp index dfa012e..411a814 100644 --- a/countseqscommand.cpp +++ b/countseqscommand.cpp @@ -8,7 +8,6 @@ */ #include "countseqscommand.h" -#include "groupmap.h" #include "sharedutilities.h" #include "counttable.h" @@ -17,6 +16,7 @@ vector CountSeqsCommand::setParameters(){ try { CommandParameter pname("name", "InputTypes", "", "", "none", "none", "none","count",false,true,true); parameters.push_back(pname); CommandParameter pgroup("group", "InputTypes", "", "", "none", "none", "none","",false,false,true); parameters.push_back(pgroup); + CommandParameter pprocessors("processors", "Number", "", "1", "", "", "","",false,false,true); parameters.push_back(pprocessors); CommandParameter plarge("large", "Boolean", "", "F", "", "", "","",false,false); parameters.push_back(plarge); CommandParameter pgroups("groups", "String", "", "", "", "", "","",false,false); parameters.push_back(pgroups); CommandParameter pinputdir("inputdir", "String", "", "", "", "", "","",false,false); parameters.push_back(pinputdir); @@ -39,6 +39,7 @@ string CountSeqsCommand::getHelpString(){ helpString += "The groups parameter allows you to indicate which groups you want to include in the counts, by default all groups in your groupfile are used.\n"; helpString += "The large parameter indicates the name and group files are too large to fit in RAM.\n"; helpString += "When you use the groups parameter and a sequence does not represent any sequences from the groups you specify it is not included in the .count.summary file.\n"; + helpString += "The processors parameter allows you to specify the number of processors to use. The default is 1.\n"; helpString += "The count.seqs command should be in the following format: count.seqs(name=yourNameFile).\n"; helpString += "Example count.seqs(name=amazon.names) or make.table(name=amazon.names).\n"; helpString += "Note: No spaces between parameter labels (i.e. name), '=' and parameters (i.e.yourNameFile).\n"; @@ -147,6 +148,10 @@ CountSeqsCommand::CountSeqsCommand(string option) { string temp = validParameter.validFile(parameters, "large", false); if (temp == "not found") { temp = "F"; } large = m->isTrue(temp); + + temp = validParameter.validFile(parameters, "processors", false); if (temp == "not found"){ temp = m->getProcessors(); } + m->setProcessors(temp); + m->mothurConvert(temp, processors); //if the user changes the output directory command factory will send this info to us in the output parameter outputDir = validParameter.validFile(parameters, "outputdir", false); if (outputDir == "not found"){ outputDir = m->hasPath(namefile); } @@ -171,11 +176,16 @@ int CountSeqsCommand::execute(){ string outputFileName = getOutputFileName("count", variables); int total = 0; + int start = time(NULL); if (!large) { total = processSmall(outputFileName); } else { total = processLarge(outputFileName); } if (m->control_pressed) { m->mothurRemove(outputFileName); return 0; } + m->mothurOut("It took " + toString(time(NULL) - start) + " secs to create a table for " + toString(total) + " sequences."); + m->mothurOutEndLine(); + m->mothurOutEndLine(); + //set rabund file as new current rabundfile itTypes = outputTypes.find("count"); if (itTypes != outputTypes.end()) { @@ -225,13 +235,172 @@ int CountSeqsCommand::processSmall(string outputFileName){ } } out << endl; + out.close(); + + int total = createProcesses(groupMap, outputFileName); + + if (groupfile != "") { delete groupMap; } + + return total; + } + catch(exception& e) { + m->errorOut(e, "CountSeqsCommand", "processSmall"); + exit(1); + } +} +/**************************************************************************************************/ +int CountSeqsCommand::createProcesses(GroupMap*& groupMap, string outputFileName) { + try { - //open input file - ifstream in; + vector processIDS; + int process = 0; + vector positions; + vector lines; + int numSeqs = 0; + +#if defined (__APPLE__) || (__MACH__) || (linux) || (__linux) || (__linux__) || (__unix__) || (__unix) + positions = m->divideFilePerLine(namefile, processors); + for (int i = 0; i < (positions.size()-1); i++) { lines.push_back(linePair(positions[i], positions[(i+1)])); } +#else + if(processors == 1){ lines.push_back(linePair(0, 1000)); } + else { + int numSeqs = 0; + positions = m->setFilePosEachLine(namefile, numSeqs); + if (positions.size() < processors) { processors = positions.size(); } + + //figure out how many sequences you have to process + int numSeqsPerProcessor = numSeqs / processors; + for (int i = 0; i < processors; i++) { + int startIndex = i * numSeqsPerProcessor; + if(i == (processors - 1)){ numSeqsPerProcessor = numSeqs - i * numSeqsPerProcessor; } + lines.push_back(linePair(positions[startIndex], numSeqsPerProcessor)); + } + } +#endif + + +#if defined (__APPLE__) || (__MACH__) || (linux) || (__linux) || (__linux__) || (__unix__) || (__unix) + + //loop through and create all the processes you want + while (process != processors-1) { + int pid = fork(); + + if (pid > 0) { + processIDS.push_back(pid); //create map from line number to pid so you can append files in correct order later + process++; + }else if (pid == 0){ + string filename = toString(getpid()) + ".temp"; + numSeqs = driver(lines[process].start, lines[process].end, filename, groupMap); + + string tempFile = toString(getpid()) + ".num.temp"; + ofstream outTemp; + m->openOutputFile(tempFile, outTemp); + + outTemp << numSeqs << endl; + outTemp.close(); + + exit(0); + }else { + m->mothurOut("[ERROR]: unable to spawn the necessary processes."); m->mothurOutEndLine(); + for (int i = 0; i < processIDS.size(); i++) { kill (processIDS[i], SIGINT); } + exit(0); + } + } + + string filename = toString(getpid()) + ".temp"; + numSeqs = driver(lines[processors-1].start, lines[processors-1].end, filename, groupMap); + + //force parent to wait until all the processes are done + for (int i=0;iopenInputFile(tempFile, intemp); + + int num; + intemp >> num; intemp.close(); + numSeqs += num; + m->mothurRemove(tempFile); + } +#else + vector pDataArray; + DWORD dwThreadIdArray[processors-1]; + HANDLE hThreadArray[processors-1]; + vector copies; + + //Create processor worker threads. + for( int i=0; igetCopy(groupMap); + copies.push_back(copyGroup); + vector cGroups = Groups; + + countData* temp = new countData(filename, copyGroup, m, lines[i].start, lines[i].end, groupfile, namefile, cGroups); + pDataArray.push_back(temp); + processIDS.push_back(i); + + hThreadArray[i] = CreateThread(NULL, 0, MyCountThreadFunction, pDataArray[i], 0, &dwThreadIdArray[i]); + } + + string filename = toString(processors-1) + ".temp"; + numSeqs = driver(lines[processors-1].start, lines[processors-1].end, filename, groupMap); + + //Wait until all threads have terminated. + WaitForMultipleObjects(processors-1, hThreadArray, TRUE, INFINITE); + + //Close all thread handles and free memory allocations. + for(int i=0; i < pDataArray.size(); i++){ + numSeqs += pDataArray[i]->total; + delete copies[i]; + CloseHandle(hThreadArray[i]); + delete pDataArray[i]; + } +#endif + + //append output files + for(int i=0;iappendFiles((toString(processIDS[i]) + ".temp"), outputFileName); + m->mothurRemove((toString(processIDS[i]) + ".temp")); + } + m->appendFiles(filename, outputFileName); + m->mothurRemove(filename); + + + //sanity check + if (groupfile != "") { + if (numSeqs != groupMap->getNumSeqs()) { + m->mothurOut("[ERROR]: processes reported processing " + toString(numSeqs) + " sequences, but group file indicates you have " + toString(groupMap->getNumSeqs()) + " sequences."); + if (processors == 1) { m->mothurOut(" Could you have a file mismatch?\n"); } + else { m->mothurOut(" Either you have a file mismatch or a process failed to complete the task assigned to it.\n"); m->control_pressed = true; } + } + } + return numSeqs; + } + catch(exception& e) { + m->errorOut(e, "CountSeqsCommand", "createProcesses"); + exit(1); + } +} +/**************************************************************************************************/ +int CountSeqsCommand::driver(unsigned long long start, unsigned long long end, string outputFileName, GroupMap*& groupMap) { + try { + + ofstream out; + m->openOutputFile(outputFileName, out); + + ifstream in; m->openInputFile(namefile, in); + in.seekg(start); - int total = 0; - while (!in.eof()) { + bool done = false; + int total = 0; + while (!done) { if (m->control_pressed) { break; } string firstCol, secondCol; @@ -240,7 +409,7 @@ int CountSeqsCommand::processSmall(string outputFileName){ m->checkName(firstCol); m->checkName(secondCol); //cout << firstCol << '\t' << secondCol << endl; - + vector names; m->splitAtChar(secondCol, names, ','); @@ -278,16 +447,23 @@ int CountSeqsCommand::processSmall(string outputFileName){ } total += names.size(); + +#if defined (__APPLE__) || (__MACH__) || (linux) || (__linux) || (__linux__) || (__unix__) || (__unix) + unsigned long long pos = in.tellg(); + if ((pos == -1) || (pos >= end)) { break; } +#else + if (in.eof()) { break; } +#endif + } in.close(); out.close(); - - if (groupfile != "") { delete groupMap; } - + return total; + } catch(exception& e) { - m->errorOut(e, "CountSeqsCommand", "processSmall"); + m->errorOut(e, "CountSeqsCommand", "driver"); exit(1); } }