mirror of
https://git.phc.dm.unipi.it/3dY_0/Calcolo_Parallelo_Cluster_Steffe.git
synced 2026-10-07 02:24:51 +00:00
MPI sorter with buffers, missing SnowPlow technique
This commit is contained in:
@@ -0,0 +1,120 @@
|
||||
# Program for compiling MPI cpp programs
|
||||
CC = mpicxx
|
||||
CXX = mpicxx
|
||||
# Extra flags to give to the processor compiler
|
||||
CFLAGS = -g
|
||||
#TODO -Wall -Werror -Wextra
|
||||
|
||||
#
|
||||
SRC = main.cpp
|
||||
OBJ = $(SRC:.cpp=.o)
|
||||
NAME = merge_sort_enhanced
|
||||
|
||||
# SBATCH parameters
|
||||
JOBNAME = Distributed_Sorting
|
||||
PARTITION = production
|
||||
TIME = 12:00:00
|
||||
MEM = 3G
|
||||
NODELIST = steffe[12-14]
|
||||
CPUS_PER_TASK = 1
|
||||
NTASKS_PER_NODE = 1
|
||||
OUTPUT = ./%x.%j.out
|
||||
ERROR = ./e%x.%j.err
|
||||
|
||||
#
|
||||
.PHONY: all run detail clean fclean re
|
||||
|
||||
.o: .cpp
|
||||
$(CC) -c $(CFLAGS) $< -o $@
|
||||
|
||||
all: $(NAME)
|
||||
|
||||
$(NAME): $(OBJ)
|
||||
$(CC) $(CFLAGS) -o $@ $^
|
||||
|
||||
run: $(NAME)
|
||||
## Se modifico il Makefile con i parametri di sbatch ricreo anche il launcher.sh
|
||||
# @if ! [ -f launcher.sh ]; then \
|
||||
|
||||
@echo "#!/bin/bash" > launcher.sh; \
|
||||
echo "## sbatch is the command line interpreter for Slurm" >> launcher.sh; \
|
||||
echo "" >> launcher.sh; \
|
||||
echo "## specify the name of the job in the queueing system" >> launcher.sh; \
|
||||
echo "#SBATCH --job-name=$(JOBNAME)" >> launcher.sh; \
|
||||
echo "## specify the partition for the resource allocation. if not specified, slurm is allowed to take the default(the one with a star *)" >> launcher.sh; \
|
||||
echo "#SBATCH --partition=$(PARTITION)" >> launcher.sh; \
|
||||
echo "## format for time is days-hours:minutes:seconds, is used as time limit for the execution duration" >> launcher.sh; \
|
||||
echo "#SBATCH --time=$(TIME)" >> launcher.sh; \
|
||||
echo "## specify the real memory required per node. suffix can be K-M-G-T but if not present is MegaBytes by default" >> launcher.sh; \
|
||||
echo "#SBATCH --mem=$(MEM)" >> launcher.sh; \
|
||||
echo "## format for hosts as a range(steffe[1-4,10-15,20]), to specify hosts needed to satisfy resource requirements" >> launcher.sh; \
|
||||
echo "#SBATCH --nodelist=$(NODELIST)" >> launcher.sh; \
|
||||
echo "## to specify the number of processors per task, default is one" >> launcher.sh; \
|
||||
echo "#SBATCH --cpus-per-task=$(CPUS_PER_TASK)" >> launcher.sh; \
|
||||
echo "## to specify the number of tasks to be invoked on each node" >> launcher.sh; \
|
||||
echo "#SBATCH --ntasks-per-node=$(NTASKS_PER_NODE)" >> launcher.sh; \
|
||||
echo "## to specify the file of utput and error" >> launcher.sh; \
|
||||
echo "#SBATCH --output $(OUTPUT)" >> launcher.sh; \
|
||||
echo "#SBATCH --error $(ERROR)" >> launcher.sh; \
|
||||
echo "" >> launcher.sh; \
|
||||
echo "mpirun $(NAME) $(ARGS)" >> launcher.sh; \
|
||||
chmod +x launcher.sh; \
|
||||
echo "The 'launcher.sh' script has been created and is ready to run."; \
|
||||
|
||||
# else \
|
||||
# chmod +x launcher.sh; \
|
||||
# fi
|
||||
@echo; sbatch launcher.sh
|
||||
@echo " To see job list you can use 'squeue'."
|
||||
@echo " To cancel a job you can use 'scancel jobid'."
|
||||
|
||||
detail:
|
||||
@echo "Compiler flags and options that mpicxx would use for compiling an MPI program: "
|
||||
@mpicxx --showme:compile
|
||||
@echo
|
||||
@echo "Linker flags and options that mpicxx would use for linking an MPI program: "
|
||||
@mpicxx --showme:link
|
||||
|
||||
clean:
|
||||
## Sembra non funzionare read
|
||||
# read -p "rm: remove all files \"./$(JOBNAME).*.out\" and \"./e$(JOBNAME).*.err\"? (y/n)" choice
|
||||
# @if [ "$$choice" = "y" ]; then \
|
||||
|
||||
@echo rm -f ./$(JOBNAME).*.out
|
||||
@for file in ./$(JOBNAME).*.out; do \
|
||||
rm -f "$$file"; \
|
||||
done
|
||||
@echo rm -f ./e$(JOBNAME).*.err
|
||||
@for file in ./e$(JOBNAME).*.err; do \
|
||||
rm -f "$$file"; \
|
||||
done
|
||||
|
||||
# fi
|
||||
rm -f *~ $(OBJ)
|
||||
|
||||
fclean: clean
|
||||
rm -f ./launcher.sh;
|
||||
rm -f ./nohup.out
|
||||
rm -f /mnt/raid/tmp/SortedRun*
|
||||
rm -f /mnt/raid/tmp/*.sort
|
||||
rm -f $(NAME)
|
||||
|
||||
re: fclean all
|
||||
|
||||
|
||||
|
||||
# mpicxx *.c
|
||||
|
||||
# mpirun/mpiexec ... //will run X copies of the program in the current run-time environment, scheduling(by default) in a round-robin fashion by CPU slot.
|
||||
|
||||
# SLIDE 5 Durastante
|
||||
# The Script
|
||||
# #!/bin/bash
|
||||
# #SBATCH --job-name=dascegliere
|
||||
# #SBATCH --mem=size[unis]
|
||||
# #SBATCH -n 10
|
||||
# #SBATCH --time=12:00:00
|
||||
# #SBATCH --nodelist=lista
|
||||
# #SBATCH --partition=ports
|
||||
# #ecc..
|
||||
# mpirun ...
|
||||
Executable
BIN
Binary file not shown.
@@ -0,0 +1,65 @@
|
||||
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <stdint.h>
|
||||
#include <sys/time.h>
|
||||
|
||||
//The number of numbers that populate the file
|
||||
#define TEST_SIZE 1
|
||||
#define BUF_SIZE 2097152//16 MegaBytes
|
||||
#define BUF_SIZE_HGB 67108864//512 MegaBytes
|
||||
#define BUF_SIZE_GB 134217728//1 GigaBytes
|
||||
#define BUF_SIZE_10GB 1342177280//10 GigaBytes
|
||||
|
||||
/*
|
||||
Generate a file of given size filling it with random 8 bytes numbers.
|
||||
Return the number of microseconds needed for the generation and writing on disk.
|
||||
*/
|
||||
long long benchmark_generate_file(const char *pathname, unsigned int seed)
|
||||
{
|
||||
struct timeval start, end;
|
||||
unsigned int size = 0;
|
||||
FILE *file;
|
||||
int64_t *buffer;
|
||||
|
||||
srand(seed);
|
||||
// if (access(pathname, F_OK) == 0)
|
||||
// return -2;//File already exist
|
||||
if (size == 0)
|
||||
{
|
||||
printf("Insert a multiple of %d(%d MegaBytes) for the size of the target file:\n", BUF_SIZE_GB, BUF_SIZE_GB/131072);
|
||||
while (1)
|
||||
{
|
||||
if (scanf("%u", &size) != 1)
|
||||
{
|
||||
printf("Insert a valid size for the file:\n");
|
||||
while (getchar() != '\n');//Clear the input buffer to prevent an infinite loop
|
||||
}
|
||||
else
|
||||
break;
|
||||
}
|
||||
printf("Future file dimension: (%d * %u) Mb\n",BUF_SIZE_GB/131072, size);
|
||||
}
|
||||
buffer = (int64_t*)malloc(BUF_SIZE_GB * sizeof(int64_t));
|
||||
if (!buffer)
|
||||
return -1;//Something went wrong
|
||||
file = fopen(pathname, "wb");
|
||||
if (!file){
|
||||
free(buffer);
|
||||
return -1;//Something went wrong
|
||||
}
|
||||
gettimeofday(&start, NULL);//Timer Start
|
||||
for (unsigned int i = 0; i < size; i++)
|
||||
{
|
||||
for(unsigned int j=0; j < BUF_SIZE_GB; j++)
|
||||
{
|
||||
buffer[j] = ((int64_t)rand() << 32) | rand();
|
||||
// printf("%ld\n", buffer[j]);
|
||||
}
|
||||
fwrite(buffer, sizeof(int64_t), BUF_SIZE_GB, file);
|
||||
}
|
||||
gettimeofday(&end, NULL);//Timer Stop
|
||||
free(buffer);
|
||||
fclose(file);
|
||||
return (end.tv_sec - start.tv_sec) * 1000000LL + (end.tv_usec - start.tv_usec);
|
||||
}
|
||||
@@ -0,0 +1,20 @@
|
||||
#include <stdio.h>
|
||||
|
||||
long long benchmark_generate_file(const char *pathname, unsigned int seed);
|
||||
long long benchmark_reader_file(const char *pathname);
|
||||
int is_sorted(const char *pathname);
|
||||
void print_partial_file(const char *pathname, unsigned long long startOffset, unsigned long long endOffset);
|
||||
void print_all_file(const char *pathname);
|
||||
|
||||
int main(int argc, char* argv[]){
|
||||
|
||||
printf("scrittura file %s: time(%lld microseconds)\n","testiamolo.bin",benchmark_generate_file("testiamolo.bin",42));//TODO cambiare il seed una volta completatp
|
||||
// printf("POI\n");
|
||||
// printf("lettura file %s: time(%lld microseconds)\n","testiamolo.bin",benchmark_reader_file("testiamolo.bin"));
|
||||
// printf("POI\n");
|
||||
// printf("sono uguali? %d\n",is_sorted("testiamolo.bin"));
|
||||
// print_partial_file("testiamolo.bin",2,4);
|
||||
// printf("POI\n");
|
||||
// print_all_file("testiamolo.bin");
|
||||
return 0;
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <stdint.h>
|
||||
#include <sys/time.h>
|
||||
|
||||
/*
|
||||
*/
|
||||
long long benchmark_reader_file(const char *pathname)
|
||||
{
|
||||
struct timeval start, end;
|
||||
unsigned long long startOffset = 0; // Start from the first number (0-based index)
|
||||
unsigned long long endOffset = -1; // Read up to the X number (exclusive), -1 in the last one because is unsigned
|
||||
FILE *file;
|
||||
int64_t num;
|
||||
|
||||
file = fopen(pathname, "rb");
|
||||
if (!file)
|
||||
return -1;//Something went wrong
|
||||
fseek(file, startOffset * sizeof(int64_t), SEEK_SET);
|
||||
gettimeofday(&start, NULL);//Timer Start
|
||||
unsigned long long currentOffset = startOffset;
|
||||
while (currentOffset < endOffset && fread(&num, sizeof(int64_t), 1, file) == 1)
|
||||
currentOffset++;
|
||||
gettimeofday(&end, NULL);//Timer Stop
|
||||
fclose(file);
|
||||
return (end.tv_sec - start.tv_sec) * 1000000LL + (end.tv_usec - start.tv_usec);
|
||||
}
|
||||
|
||||
/*
|
||||
*/
|
||||
int is_sorted(const char *pathname)
|
||||
{
|
||||
unsigned long long startOffset = 0; // Start from the first number (0-based index)
|
||||
unsigned long long endOffset = -1; // Read up to the X number (exclusive), -1 in the last one because is unsigned
|
||||
FILE *file;
|
||||
int64_t num;
|
||||
long long int count=1;
|
||||
|
||||
file = fopen(pathname, "rb");
|
||||
if (!file)
|
||||
return -1;//Something went wrong
|
||||
fseek(file, startOffset * sizeof(int64_t), SEEK_SET);
|
||||
unsigned long long currentOffset = startOffset;
|
||||
int64_t tmp;
|
||||
fread(&tmp, sizeof(int64_t), 1, file); // Take first element(number) in the file
|
||||
currentOffset++;
|
||||
while (currentOffset < endOffset && fread(&num, sizeof(int64_t), 1, file) == 1)
|
||||
{
|
||||
count++;
|
||||
if(tmp > num){
|
||||
fclose(file);
|
||||
// printf("non ordinati\n");
|
||||
return 0;
|
||||
}
|
||||
tmp = num;
|
||||
currentOffset++;
|
||||
}
|
||||
fclose(file);
|
||||
// printf("%lld numeri ordinati",count);
|
||||
return 1;
|
||||
}
|
||||
|
||||
/*
|
||||
Start from the start number (0-based index), read up to the end number (exclusive).
|
||||
*/
|
||||
void print_partial_file(const char *pathname, unsigned long long startOffset, unsigned long long endOffset)
|
||||
{
|
||||
FILE *file;
|
||||
int64_t num;
|
||||
|
||||
file = fopen(pathname, "rb");
|
||||
if (!file)
|
||||
return;//Something went wrong
|
||||
fseek(file, startOffset * sizeof(int64_t), SEEK_SET);
|
||||
unsigned long long currentOffset = startOffset;
|
||||
while (currentOffset < endOffset && fread(&num, sizeof(int64_t), 1, file) == 1)
|
||||
{
|
||||
printf("%lld\n", (long long)num);
|
||||
currentOffset++;
|
||||
}
|
||||
fclose(file);
|
||||
return;
|
||||
}
|
||||
|
||||
/*
|
||||
*/
|
||||
void print_all_file(const char *pathname)
|
||||
{
|
||||
unsigned long long startOffset = 0; // Start from the first number (0-based index)
|
||||
unsigned long long endOffset = -1; // Read up to the X number (exclusive), -1 in the last one because is unsigned
|
||||
FILE *file;
|
||||
int64_t num;
|
||||
|
||||
file = fopen(pathname, "rb");
|
||||
if (!file)
|
||||
return;//Something went wrong
|
||||
fseek(file, startOffset * sizeof(int64_t), SEEK_SET);
|
||||
unsigned long long currentOffset = startOffset;
|
||||
while (currentOffset < endOffset && fread(&num, sizeof(int64_t), 1, file) == 1)
|
||||
{
|
||||
printf("%lld - ", (long long)num);
|
||||
currentOffset++;
|
||||
}
|
||||
printf("EOF\n");
|
||||
fclose(file);
|
||||
return;
|
||||
}
|
||||
@@ -0,0 +1,357 @@
|
||||
#include <mpi.h>
|
||||
#include <fcntl.h>
|
||||
#include <dirent.h>
|
||||
#include <unistd.h>
|
||||
#include <sys/stat.h>
|
||||
#include <sys/types.h>
|
||||
#include <queue>
|
||||
#include <ctime>
|
||||
#include <string>
|
||||
#include <vector>
|
||||
#include <cstdio>
|
||||
#include <fstream>
|
||||
#include <cstdlib>
|
||||
#include <iostream>
|
||||
#include <algorithm>
|
||||
|
||||
/*
|
||||
8-Byte numbers in 256KB = 32768
|
||||
8-Byte numbers in 1MB = 131072
|
||||
8-Byte numbers in 1GB = 134217728
|
||||
|
||||
All the programm assume numbers as 64_bits.
|
||||
To visualize binary files in bash can be used:
|
||||
od -t d8 -A n binaryfile.bin #For in use format
|
||||
od -t d8 -A n --endian=little binaryfile.bin #For little-endian format
|
||||
od -t d8 -A n --endian=big binaryfile.bin #For big-endian format
|
||||
*/
|
||||
#define BUFFERSIZE 32768
|
||||
#define CACHENUM 130000
|
||||
#define RAMNUM 268435456
|
||||
#define ALLOW_BUFFER 1
|
||||
#define ALLOW_SNOWPLOW 1
|
||||
|
||||
void sortRuns(unsigned long long fileSize, unsigned long long sliceSize, unsigned long long maxLoop, FILE* file, int id, int mpiRank, int mpiSize)
|
||||
{
|
||||
unsigned long long startOffset, endOffset, currentOffset; //The interval is [startOffset, endOffset)
|
||||
double startTot, start, end;
|
||||
int64_t num;
|
||||
std::vector<int64_t> bigVect;
|
||||
int64_t buffer[BUFFERSIZE];
|
||||
bigVect.reserve(sliceSize);
|
||||
|
||||
startTot = MPI_Wtime(); //Microsecond precision. Can't use time(), because each process will have a different "zero" time
|
||||
start = MPI_Wtime();
|
||||
for(unsigned long long l = 0; l < maxLoop; l++) //Populate the vector with the values in the file
|
||||
{
|
||||
startOffset = sliceSize * (mpiRank + (mpiSize * l));
|
||||
if (startOffset >= fileSize)
|
||||
break;
|
||||
endOffset = startOffset + sliceSize;
|
||||
fseek(file, startOffset * sizeof(int64_t), SEEK_SET);
|
||||
currentOffset = startOffset;
|
||||
bigVect.clear();
|
||||
|
||||
if (ALLOW_BUFFER) //Branch to test performance with and without buffer
|
||||
{
|
||||
while (currentOffset < endOffset)
|
||||
{
|
||||
unsigned long long elementsToRead = std::min(endOffset - currentOffset, static_cast<unsigned long long>(BUFFERSIZE)); //It's important to check because if the difference between endOffset and startOffset is smaller than BUFFERSIZE we don't have to read further
|
||||
unsigned long long elementsRead = fread(buffer, sizeof(int64_t), elementsToRead, file);
|
||||
|
||||
for (unsigned long long i = 0; i < elementsRead; ++i)
|
||||
{
|
||||
bigVect.push_back(buffer[i]);
|
||||
}
|
||||
currentOffset += elementsRead; //Increment currentOffset based on the number of elements read
|
||||
if (elementsRead < BUFFERSIZE) // Check if we have reached the end of the file
|
||||
break;
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
while (currentOffset < endOffset && fread(&num, sizeof(int64_t), 1, file) == 1)
|
||||
{
|
||||
bigVect.push_back(num);
|
||||
currentOffset++;
|
||||
}
|
||||
}
|
||||
end = MPI_Wtime();
|
||||
std::cout << " " << end-start << "s" << " => Time to read file from offset " << startOffset << " to " << endOffset << " in Process " << mpiRank+1 << "/" << mpiSize << " memory" << std::endl;
|
||||
start = MPI_Wtime();
|
||||
|
||||
sort(bigVect.begin(), bigVect.end());
|
||||
|
||||
end = MPI_Wtime();
|
||||
std::cout << " " << end-start << "s" << " => Time to sort elements in Process " << mpiRank+1 << "/" << mpiSize << " memory" << std::endl;
|
||||
start = MPI_Wtime();
|
||||
|
||||
std::string templateName = "/mnt/raid/tmp/SortedRun" + std::to_string(id) + "_XXXXXX"; //If absolute path does not exist the temporary file will not be created
|
||||
int tmpFile = mkstemp(&templateName[0]); //Create a temporary file based on template
|
||||
if (tmpFile == -1)
|
||||
{
|
||||
std::cout << "Error creating temporary file" << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
for (unsigned long long i = 0; i < bigVect.size(); ++i) //Write the ordered number in a temp file
|
||||
{
|
||||
if (ALLOW_BUFFER) //Branch to test performance with and without buffer
|
||||
{
|
||||
buffer[i % BUFFERSIZE] = bigVect[i];
|
||||
if ((i + 1) % BUFFERSIZE == 0 || i == bigVect.size() - 1)
|
||||
{
|
||||
ssize_t tw = write(tmpFile, buffer, sizeof(int64_t) * ((i % BUFFERSIZE) + 1));
|
||||
if (tw == -1)
|
||||
{
|
||||
std::cout << "Error writing to file" << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
int64_t elem = bigVect[i];
|
||||
ssize_t tw = write(tmpFile, &elem, sizeof(int64_t));
|
||||
if (tw == -1)
|
||||
{
|
||||
std::cout << "Error writing to file" << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
}
|
||||
}
|
||||
off_t sz = lseek(tmpFile, 0, SEEK_END);
|
||||
if (sz == 0)
|
||||
{
|
||||
if (close(tmpFile) == -1)
|
||||
{
|
||||
std::cout << "Error closing the file" << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
if (unlink(&templateName[0]) == -1)
|
||||
{
|
||||
std::cout << "Error unlinking file" << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
}
|
||||
if (close(tmpFile) == -1)
|
||||
{
|
||||
std::cout << "Error closing the file" << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
end = MPI_Wtime();
|
||||
std::cout << " " << end-start << "s" << " => Time to write '" << templateName << "' and fill it up with " << sz/8 << " sorted elements by Process " << mpiRank+1 << "/" << mpiSize << std::endl;
|
||||
start = MPI_Wtime();
|
||||
}
|
||||
end = MPI_Wtime();
|
||||
std::cout << end-startTot << "s" << " => Time function sortRuns() in Process " << mpiRank+1 << "/" << mpiSize << std::endl;
|
||||
}
|
||||
|
||||
void snowPlowRuns(unsigned long long fileSize, unsigned long long sliceSize, unsigned long long maxLoop, FILE* file, int id, int mpiRank, int mpiSize)
|
||||
{
|
||||
if (ALLOW_SNOWPLOW)
|
||||
std::cout << "Can't compute files of size bigger then " << RAMNUM * mpiSize / 134217728 << "Gb with " << mpiSize << " processes (currently file is " << fileSize / 134217728 << "Gb)" << std::endl;
|
||||
else
|
||||
{
|
||||
maxLoop = (fileSize / (RAMNUM * mpiSize)) + 1;
|
||||
sortRuns(fileSize, RAMNUM, maxLoop, file, id, mpiRank, mpiSize);
|
||||
}
|
||||
}
|
||||
|
||||
void kMerge(const std::string &argFile, int id, int mpiRank, int mpiSize)
|
||||
{
|
||||
std::string fileDir = "/mnt/raid/tmp/";
|
||||
std::string pattern = "SortedRun" + std::to_string(id) + "_";
|
||||
std::vector<int> fds; //To store the file descriptor of each file to merge
|
||||
std::vector<std::string> fns; //To store the file name of each file to delete after merge
|
||||
size_t lastSlash = argFile.find_last_of('/');
|
||||
std::string nameOnly = (lastSlash != std::string::npos) ? argFile.substr(lastSlash + 1) : argFile;
|
||||
std::string finalFile = "/mnt/raid/tmp/" + nameOnly + (ALLOW_BUFFER == 1 ? ".buf" : "") + ".sort";
|
||||
double start, end;
|
||||
int fileCount = 0;
|
||||
|
||||
DIR *dir = opendir(fileDir.c_str());
|
||||
if (dir)
|
||||
{
|
||||
struct dirent *entry;
|
||||
while ((entry = readdir(dir)) != nullptr)
|
||||
{
|
||||
if (entry->d_type == DT_REG) //Check if it's a regular file
|
||||
{
|
||||
std::string filename = entry->d_name;
|
||||
if (filename.find(pattern) != std::string::npos) //Check if the file name matches the pattern
|
||||
{
|
||||
std::string tmpFile = fileDir + "/" + filename;
|
||||
int fd = open(tmpFile.c_str(), O_RDONLY); //Open the file and save the file descriptor
|
||||
if (fd != -1)
|
||||
{
|
||||
fds.push_back(fd);
|
||||
fns.push_back(tmpFile);
|
||||
fileCount++;
|
||||
}
|
||||
else
|
||||
std::cout << "Error opening file '" << tmpFile << "' by Process " << mpiRank+1 << "/" << mpiSize << std::endl;
|
||||
}
|
||||
}
|
||||
}
|
||||
closedir(dir);
|
||||
}
|
||||
else
|
||||
{
|
||||
std::cout << "Error opening directory '" << fileDir << "' by Process " << mpiRank+1 << "/" << mpiSize << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
|
||||
int fdFinal = open(finalFile.c_str(), O_WRONLY | O_CREAT | O_EXCL, S_IRUSR | S_IWUSR); //Open the file for writing only, creating it if it doesn't exist and not overwrite if it exists
|
||||
if (fdFinal == -1)
|
||||
{
|
||||
std::cout << "Error opening or creating final file '" << finalFile << "' by Process " << mpiRank+1 << "/" << mpiSize << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
|
||||
std::cout << std::endl << "Starting the merge process for " << fileCount << " files" << std::endl << std::endl;
|
||||
start = MPI_Wtime();
|
||||
|
||||
std::priority_queue<std::pair<int64_t, int>, std::vector<std::pair<int64_t, int>>, std::greater<std::pair<int64_t, int>>> minHeap; //Creating a Min Heap using a priority queue
|
||||
int64_t tmpValue;
|
||||
for (int fd : fds) //Populate the Min Heap with initial values from each file descriptor
|
||||
{
|
||||
if (read(fd, &tmpValue, sizeof(int64_t)) == sizeof(int64_t))
|
||||
minHeap.push({tmpValue, fd});
|
||||
else
|
||||
{
|
||||
std::cout << "Error reading from file descriptor by Process " << mpiRank+1 << "/" << mpiSize << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
}
|
||||
|
||||
int tmpfd;
|
||||
int64_t tmpValue2;
|
||||
int64_t buffer[BUFFERSIZE];
|
||||
unsigned long long i = 0;
|
||||
while (!minHeap.empty()) //Write sorted elements to the temporary file
|
||||
{
|
||||
tmpValue = minHeap.top().first;
|
||||
tmpfd = minHeap.top().second;
|
||||
if (read(tmpfd, &tmpValue2, sizeof(int64_t)) == sizeof(int64_t)) //Read another integer from the same file descriptor
|
||||
{
|
||||
minHeap.pop();
|
||||
minHeap.push({tmpValue2, tmpfd});
|
||||
}
|
||||
else //If no more values can be read
|
||||
{
|
||||
minHeap.pop();
|
||||
if (close(tmpfd) == -1)
|
||||
{
|
||||
std::cout << "Error closing the file descriptor by Process " << mpiRank+1 << "/" << mpiSize << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
}
|
||||
if (ALLOW_BUFFER) //Branch to test performance with and without buffer
|
||||
{
|
||||
buffer[i % BUFFERSIZE] = tmpValue;
|
||||
if ((i + 1) % BUFFERSIZE == 0 || minHeap.empty())
|
||||
{
|
||||
ssize_t tw = write(fdFinal, buffer, sizeof(int64_t) * ((i % BUFFERSIZE) + 1));
|
||||
if (tw == -1)
|
||||
{
|
||||
std::cout << "Error writing to file" << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
}
|
||||
i++;
|
||||
}
|
||||
else
|
||||
{
|
||||
ssize_t tw = write(fdFinal, &tmpValue, sizeof(int64_t));
|
||||
if (tw == -1)
|
||||
{
|
||||
std::cout << "Error writing to file" << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for (const std::string &fn : fns) //Remove all temporary files after merging them
|
||||
{
|
||||
if (unlink(&fn[0]) == -1)
|
||||
{
|
||||
std::cout << "Error unlinking file '" << fn << "' by Process " << mpiRank+1 << "/" << mpiSize << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
}
|
||||
end = MPI_Wtime();
|
||||
std::cout << end-start << "s" << " => Time function kMerge() in Process " << mpiRank+1 << "/" << mpiSize << std::endl;
|
||||
std::cout << std::endl << "Sorted file '" << finalFile << "'" << std::endl;
|
||||
}
|
||||
|
||||
int main(int argc, char* argv[])
|
||||
{
|
||||
MPI_Init(&argc, &argv); //Initialize the MPI environment
|
||||
|
||||
double startGlobal, endGlobal;
|
||||
int id, mpiSize, mpiRank;
|
||||
MPI_Comm_size(MPI_COMM_WORLD, &mpiSize); //Get the number of processes
|
||||
MPI_Comm_rank(MPI_COMM_WORLD, &mpiRank); //Get the number of process
|
||||
if (mpiRank == 0)
|
||||
{
|
||||
startGlobal = MPI_Wtime();
|
||||
std::srand(std::time(0));
|
||||
id = std::rand() % 10000; //Get a random id number to recognize files of different executions
|
||||
}
|
||||
MPI_Bcast(&id, 1, MPI_INT, 0, MPI_COMM_WORLD);
|
||||
|
||||
if (argc != 2)
|
||||
{
|
||||
if (mpiRank == 0)
|
||||
{
|
||||
std::cout << "Usage: " << argv[0] << " <file_to_parse>" << std::endl;
|
||||
std::cout << "It returns a file with extension '.sort' in the same directory of the not-parsed one. Make sure to have space before." << std::endl;
|
||||
std::cout << "Use arguments in the make as ARGS=\"stuff\". Example 'make run ARGS=\"/path/to/file\"'." << std::endl;
|
||||
}
|
||||
MPI_Finalize(); //Clean up the MPI environment
|
||||
return 0;
|
||||
}
|
||||
|
||||
FILE *file;
|
||||
unsigned long long fileSize, slices, sliceSize, maxLoop;
|
||||
file = fopen(argv[1], "rb"); //Open the file in mode rb (read binary)
|
||||
if(!file)
|
||||
{
|
||||
std::cout << "Error opening file: " << argv[1] << std::endl;
|
||||
MPI_Abort(MPI_COMM_WORLD, 1);
|
||||
}
|
||||
fseek(file,0,SEEK_END);
|
||||
fileSize = ftell(file) / 8; //Size in bytes of the file, correspond to the number of numbers to parse. Each number is 8 bytes
|
||||
if (mpiRank == 0)
|
||||
std::cout << "Sorting file '" << argv[1] << "' of " << fileSize << " elements" << std::endl << std::endl;
|
||||
|
||||
if (fileSize < (CACHENUM * mpiSize))
|
||||
slices = (fileSize / CACHENUM) + 1;
|
||||
else if (fileSize < (RAMNUM * mpiSize)) //TODO add more granularity considering double RAM for snow plow technique
|
||||
slices = (fileSize / RAMNUM) + 1;
|
||||
else
|
||||
slices = mpiSize + 1;
|
||||
sliceSize = (fileSize / slices) + 1; //Each process divides a number of 8-byte integers based on the size of the starting file, Attualmente dentro create Runs
|
||||
maxLoop = (slices / mpiSize) + 1;
|
||||
if (sliceSize > RAMNUM)
|
||||
snowPlowRuns(fileSize, sliceSize, maxLoop, file, id, mpiRank, mpiSize);
|
||||
else
|
||||
sortRuns(fileSize, sliceSize, maxLoop, file, id, mpiRank, mpiSize);
|
||||
fclose(file);
|
||||
|
||||
MPI_Barrier(MPI_COMM_WORLD); //Blocks the caller until all processes in the communicator have called it
|
||||
|
||||
if(mpiRank==0)
|
||||
{
|
||||
kMerge(argv[1], id, mpiRank, mpiSize);
|
||||
|
||||
endGlobal = MPI_Wtime();
|
||||
std::cout << (endGlobal-startGlobal)/60.0 << "min" << " => FULL EXECUTION TIME" << std::endl;
|
||||
std::cout << std::endl << "To visualize binary files in bash can be used:" << std::endl;
|
||||
std::cout << "od -t d8 -A n binaryfile.bin #For in use format" << std::endl;
|
||||
std::cout << "od -t d8 -A n --endian=little binaryfile.bin #For little-endian format" << std::endl;
|
||||
std::cout << "od -t d8 -A n --endian=big binaryfile.bin #For big-endian format" << std::endl;
|
||||
}
|
||||
MPI_Finalize(); //Clean up the MPI environment
|
||||
return 0;
|
||||
}
|
||||
Reference in New Issue
Block a user