На CPU вычисления
This commit is contained in:
69
src/main.cpp
69
src/main.cpp
@@ -6,12 +6,15 @@
|
||||
#include "csv_loader.hpp"
|
||||
#include "utils.hpp"
|
||||
#include "record.hpp"
|
||||
#include "day_stats.hpp"
|
||||
#include "aggregation.hpp"
|
||||
#include "intervals.hpp"
|
||||
#include "gpu_loader.hpp"
|
||||
|
||||
// Функция: отобрать записи для конкретного ранга
|
||||
std::vector<Record> select_records_for_rank(
|
||||
const std::map<long long, std::vector<Record>>& days,
|
||||
const std::vector<long long>& day_list)
|
||||
const std::map<DayIndex, std::vector<Record>>& days,
|
||||
const std::vector<DayIndex>& day_list)
|
||||
{
|
||||
std::vector<Record> out;
|
||||
for (auto d : day_list) {
|
||||
@@ -52,7 +55,7 @@ int main(int argc, char** argv) {
|
||||
continue;
|
||||
}
|
||||
|
||||
int count = vec.size();
|
||||
int count = static_cast<int>(vec.size());
|
||||
MPI_Send(&count, 1, MPI_INT, r, 0, MPI_COMM_WORLD);
|
||||
MPI_Send(vec.data(), count * sizeof(Record), MPI_BYTE, r, 1, MPI_COMM_WORLD);
|
||||
}
|
||||
@@ -72,14 +75,72 @@ int main(int argc, char** argv) {
|
||||
std::cout << "Rank " << rank << " received "
|
||||
<< local_records.size() << " records" << std::endl;
|
||||
|
||||
// ====== АГРЕГАЦИЯ НА КАЖДОМ УЗЛЕ ======
|
||||
auto local_stats = aggregate_days(local_records);
|
||||
std::cout << "Rank " << rank << " aggregated "
|
||||
<< local_stats.size() << " days" << std::endl;
|
||||
|
||||
// ====== СБОР АГРЕГИРОВАННЫХ ДАННЫХ НА RANK 0 ======
|
||||
std::vector<DayStats> all_stats;
|
||||
|
||||
if (rank == 0) {
|
||||
// Добавляем свои данные
|
||||
all_stats.insert(all_stats.end(), local_stats.begin(), local_stats.end());
|
||||
|
||||
// Получаем данные от других узлов
|
||||
for (int r = 1; r < size; r++) {
|
||||
int count = 0;
|
||||
MPI_Recv(&count, 1, MPI_INT, r, 2, MPI_COMM_WORLD, MPI_STATUS_IGNORE);
|
||||
|
||||
std::vector<DayStats> remote_stats(count);
|
||||
MPI_Recv(remote_stats.data(), count * sizeof(DayStats),
|
||||
MPI_BYTE, r, 3, MPI_COMM_WORLD, MPI_STATUS_IGNORE);
|
||||
|
||||
all_stats.insert(all_stats.end(), remote_stats.begin(), remote_stats.end());
|
||||
}
|
||||
} else {
|
||||
// Отправляем свои агрегированные данные на rank 0
|
||||
int count = static_cast<int>(local_stats.size());
|
||||
MPI_Send(&count, 1, MPI_INT, 0, 2, MPI_COMM_WORLD);
|
||||
MPI_Send(local_stats.data(), count * sizeof(DayStats), MPI_BYTE, 0, 3, MPI_COMM_WORLD);
|
||||
}
|
||||
|
||||
// ====== ВЫЧИСЛЕНИЕ ИНТЕРВАЛОВ НА RANK 0 ======
|
||||
if (rank == 0) {
|
||||
std::cout << "Rank 0: merging " << all_stats.size() << " day stats..." << std::endl;
|
||||
|
||||
// Объединяем и сортируем
|
||||
auto merged_stats = merge_day_stats(all_stats);
|
||||
std::cout << "Rank 0: total " << merged_stats.size() << " unique days" << std::endl;
|
||||
|
||||
// Вычисляем интервалы
|
||||
auto intervals = find_intervals(merged_stats, 0.10);
|
||||
std::cout << "Found " << intervals.size() << " intervals with >=10% change" << std::endl;
|
||||
|
||||
// Записываем результат
|
||||
write_intervals("../result.csv", intervals);
|
||||
std::cout << "Results written to result.csv" << std::endl;
|
||||
|
||||
// Выводим первые несколько интервалов
|
||||
std::cout << "\nFirst 5 intervals:\n";
|
||||
std::cout << "start_date,end_date,min_open,max_close,change\n";
|
||||
for (size_t i = 0; i < std::min(intervals.size(), size_t(5)); i++) {
|
||||
const auto& iv = intervals[i];
|
||||
std::cout << day_index_to_date(iv.start_day) << ","
|
||||
<< day_index_to_date(iv.end_day) << ","
|
||||
<< iv.min_open << ","
|
||||
<< iv.max_close << ","
|
||||
<< iv.change << "\n";
|
||||
}
|
||||
}
|
||||
|
||||
// Проверка GPU (оставляем как есть)
|
||||
auto gpu_is_available = load_gpu_is_available();
|
||||
int have_gpu = 0;
|
||||
if (gpu_is_available) {
|
||||
std::cout << "Rank " << rank << " dll loaded" << std::endl;
|
||||
have_gpu = gpu_is_available();
|
||||
}
|
||||
|
||||
std::cout << "Rank " << rank << ": gpu_available=" << have_gpu << "\n";
|
||||
|
||||
MPI_Finalize();
|
||||
|
||||
Reference in New Issue
Block a user