Агрегация
This commit is contained in:
@@ -1,87 +1,48 @@
|
|||||||
#include "aggregation.hpp"
|
#include "aggregation.hpp"
|
||||||
|
#include <map>
|
||||||
#include <algorithm>
|
#include <algorithm>
|
||||||
#include <limits>
|
#include <limits>
|
||||||
#include <cmath>
|
|
||||||
|
|
||||||
std::vector<DayStats> aggregate_days(const std::vector<Record>& records) {
|
std::vector<DayStats> aggregate_days(const std::vector<Record>& records) {
|
||||||
// Группируем записи по дням
|
// Накопители для каждого дня
|
||||||
std::map<DayIndex, std::vector<const Record*>> day_records;
|
struct DayAccumulator {
|
||||||
|
double avg_sum = 0.0;
|
||||||
|
double open_min = std::numeric_limits<double>::max();
|
||||||
|
double open_max = std::numeric_limits<double>::lowest();
|
||||||
|
double close_min = std::numeric_limits<double>::max();
|
||||||
|
double close_max = std::numeric_limits<double>::lowest();
|
||||||
|
int64_t count = 0;
|
||||||
|
};
|
||||||
|
|
||||||
|
std::map<DayIndex, DayAccumulator> days;
|
||||||
|
|
||||||
for (const auto& r : records) {
|
for (const auto& r : records) {
|
||||||
DayIndex day = static_cast<DayIndex>(r.timestamp) / 86400;
|
DayIndex day = static_cast<DayIndex>(r.timestamp) / 86400;
|
||||||
day_records[day].push_back(&r);
|
auto& acc = days[day];
|
||||||
|
|
||||||
|
double avg = (r.low + r.high) / 2.0;
|
||||||
|
acc.avg_sum += avg;
|
||||||
|
acc.open_min = std::min(acc.open_min, r.open);
|
||||||
|
acc.open_max = std::max(acc.open_max, r.open);
|
||||||
|
acc.close_min = std::min(acc.close_min, r.close);
|
||||||
|
acc.close_max = std::max(acc.close_max, r.close);
|
||||||
|
acc.count++;
|
||||||
}
|
}
|
||||||
|
|
||||||
std::vector<DayStats> result;
|
std::vector<DayStats> result;
|
||||||
result.reserve(day_records.size());
|
result.reserve(days.size());
|
||||||
|
|
||||||
for (auto& [day, recs] : day_records) {
|
for (const auto& [day, acc] : days) {
|
||||||
// Сортируем по timestamp для определения first/last
|
|
||||||
std::sort(recs.begin(), recs.end(),
|
|
||||||
[](const Record* a, const Record* b) {
|
|
||||||
return a->timestamp < b->timestamp;
|
|
||||||
});
|
|
||||||
|
|
||||||
DayStats stats;
|
DayStats stats;
|
||||||
stats.day = day;
|
stats.day = day;
|
||||||
stats.low = std::numeric_limits<double>::max();
|
stats.avg = acc.avg_sum / static_cast<double>(acc.count);
|
||||||
stats.high = std::numeric_limits<double>::lowest();
|
stats.open_min = acc.open_min;
|
||||||
stats.open = recs.front()->open;
|
stats.open_max = acc.open_max;
|
||||||
stats.close = recs.back()->close;
|
stats.close_min = acc.close_min;
|
||||||
stats.first_ts = recs.front()->timestamp;
|
stats.close_max = acc.close_max;
|
||||||
stats.last_ts = recs.back()->timestamp;
|
stats.count = acc.count;
|
||||||
|
|
||||||
for (const auto* r : recs) {
|
|
||||||
stats.low = std::min(stats.low, r->low);
|
|
||||||
stats.high = std::max(stats.high, r->high);
|
|
||||||
}
|
|
||||||
|
|
||||||
stats.avg = (stats.low + stats.high) / 2.0;
|
|
||||||
|
|
||||||
result.push_back(stats);
|
result.push_back(stats);
|
||||||
}
|
}
|
||||||
|
|
||||||
return result;
|
return result;
|
||||||
}
|
}
|
||||||
|
|
||||||
std::vector<DayStats> merge_day_stats(const std::vector<DayStats>& all_stats) {
|
|
||||||
// Объединяем статистику по одинаковым дням (если такие есть)
|
|
||||||
std::map<DayIndex, DayStats> merged;
|
|
||||||
|
|
||||||
for (const auto& s : all_stats) {
|
|
||||||
auto it = merged.find(s.day);
|
|
||||||
if (it == merged.end()) {
|
|
||||||
merged[s.day] = s;
|
|
||||||
} else {
|
|
||||||
// Объединяем данные за один день
|
|
||||||
auto& m = it->second;
|
|
||||||
m.low = std::min(m.low, s.low);
|
|
||||||
m.high = std::max(m.high, s.high);
|
|
||||||
|
|
||||||
// open берём от записи с меньшим timestamp
|
|
||||||
if (s.first_ts < m.first_ts) {
|
|
||||||
m.open = s.open;
|
|
||||||
m.first_ts = s.first_ts;
|
|
||||||
}
|
|
||||||
|
|
||||||
// close берём от записи с большим timestamp
|
|
||||||
if (s.last_ts > m.last_ts) {
|
|
||||||
m.close = s.close;
|
|
||||||
m.last_ts = s.last_ts;
|
|
||||||
}
|
|
||||||
|
|
||||||
m.avg = (m.low + m.high) / 2.0;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Преобразуем в отсортированный вектор
|
|
||||||
std::vector<DayStats> result;
|
|
||||||
result.reserve(merged.size());
|
|
||||||
|
|
||||||
for (auto& [day, stats] : merged) {
|
|
||||||
result.push_back(stats);
|
|
||||||
}
|
|
||||||
|
|
||||||
return result;
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|||||||
@@ -3,12 +3,6 @@
|
|||||||
#include "record.hpp"
|
#include "record.hpp"
|
||||||
#include "day_stats.hpp"
|
#include "day_stats.hpp"
|
||||||
#include <vector>
|
#include <vector>
|
||||||
#include <map>
|
|
||||||
|
|
||||||
// Агрегация записей по дням на одном узле
|
// Агрегация записей по дням на одном узле
|
||||||
std::vector<DayStats> aggregate_days(const std::vector<Record>& records);
|
std::vector<DayStats> aggregate_days(const std::vector<Record>& records);
|
||||||
|
|
||||||
// Объединение агрегированных данных с разных узлов
|
|
||||||
// (на случай если один день попал на разные узлы - но в нашей схеме это не должно случиться)
|
|
||||||
std::vector<DayStats> merge_day_stats(const std::vector<DayStats>& all_stats);
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,28 +1,15 @@
|
|||||||
#pragma once
|
#pragma once
|
||||||
#include <cstdint>
|
#include <cstdint>
|
||||||
|
|
||||||
using DayIndex = long long;
|
using DayIndex = int64_t;
|
||||||
|
|
||||||
// Агрегированные данные за один день
|
// Агрегированные данные за один день
|
||||||
struct DayStats {
|
struct DayStats {
|
||||||
DayIndex day; // индекс дня (timestamp / 86400)
|
DayIndex day; // индекс дня (timestamp / 86400)
|
||||||
double low; // минимальный Low за день
|
double avg; // среднее значение (Low + High) / 2 по всем записям
|
||||||
double high; // максимальный High за день
|
double open_min; // минимальный Open за день
|
||||||
double open; // первый Open за день
|
double open_max; // максимальный Open за день
|
||||||
double close; // последний Close за день
|
double close_min; // минимальный Close за день
|
||||||
double avg; // среднее = (low + high) / 2
|
double close_max; // максимальный Close за день
|
||||||
double first_ts; // timestamp первой записи (для определения порядка open)
|
int64_t count; // количество записей, по которым агрегировали
|
||||||
double last_ts; // timestamp последней записи (для определения close)
|
|
||||||
};
|
};
|
||||||
|
|
||||||
// Интервал с изменением >= 10%
|
|
||||||
struct Interval {
|
|
||||||
DayIndex start_day;
|
|
||||||
DayIndex end_day;
|
|
||||||
double min_open;
|
|
||||||
double max_close;
|
|
||||||
double start_avg;
|
|
||||||
double end_avg;
|
|
||||||
double change;
|
|
||||||
};
|
|
||||||
|
|
||||||
|
|||||||
@@ -127,13 +127,12 @@ bool aggregate_days_gpu(
|
|||||||
for (const auto& gs : gpu_stats) {
|
for (const auto& gs : gpu_stats) {
|
||||||
DayStats ds;
|
DayStats ds;
|
||||||
ds.day = gs.day;
|
ds.day = gs.day;
|
||||||
ds.low = gs.low;
|
|
||||||
ds.high = gs.high;
|
|
||||||
ds.open = gs.open;
|
|
||||||
ds.close = gs.close;
|
|
||||||
ds.avg = gs.avg;
|
ds.avg = gs.avg;
|
||||||
ds.first_ts = gs.first_ts;
|
ds.open_min = gs.open_min;
|
||||||
ds.last_ts = gs.last_ts;
|
ds.open_max = gs.open_max;
|
||||||
|
ds.close_min = gs.close_min;
|
||||||
|
ds.close_max = gs.close_max;
|
||||||
|
ds.count = gs.count;
|
||||||
out_stats.push_back(ds);
|
out_stats.push_back(ds);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -20,13 +20,12 @@ struct GpuRecord {
|
|||||||
|
|
||||||
struct GpuDayStats {
|
struct GpuDayStats {
|
||||||
long long day;
|
long long day;
|
||||||
double low;
|
|
||||||
double high;
|
|
||||||
double open;
|
|
||||||
double close;
|
|
||||||
double avg;
|
double avg;
|
||||||
double first_ts;
|
double open_min;
|
||||||
double last_ts;
|
double open_max;
|
||||||
|
double close_min;
|
||||||
|
double close_max;
|
||||||
|
long long count;
|
||||||
};
|
};
|
||||||
|
|
||||||
using gpu_aggregate_days_fn = int (*)(
|
using gpu_aggregate_days_fn = int (*)(
|
||||||
|
|||||||
@@ -23,13 +23,12 @@ struct GpuRecord {
|
|||||||
|
|
||||||
struct GpuDayStats {
|
struct GpuDayStats {
|
||||||
long long day;
|
long long day;
|
||||||
double low;
|
|
||||||
double high;
|
|
||||||
double open;
|
|
||||||
double close;
|
|
||||||
double avg;
|
double avg;
|
||||||
double first_ts;
|
double open_min;
|
||||||
double last_ts;
|
double open_max;
|
||||||
|
double close_min;
|
||||||
|
double close_max;
|
||||||
|
long long count;
|
||||||
};
|
};
|
||||||
|
|
||||||
extern "C" int gpu_is_available() {
|
extern "C" int gpu_is_available() {
|
||||||
@@ -63,32 +62,30 @@ __global__ void aggregate_kernel(
|
|||||||
|
|
||||||
GpuDayStats stats;
|
GpuDayStats stats;
|
||||||
stats.day = day_indices[d];
|
stats.day = day_indices[d];
|
||||||
stats.low = DBL_MAX;
|
stats.open_min = DBL_MAX;
|
||||||
stats.high = -DBL_MAX;
|
stats.open_max = -DBL_MAX;
|
||||||
stats.first_ts = DBL_MAX;
|
stats.close_min = DBL_MAX;
|
||||||
stats.last_ts = -DBL_MAX;
|
stats.close_max = -DBL_MAX;
|
||||||
stats.open = 0;
|
stats.count = count;
|
||||||
stats.close = 0;
|
|
||||||
|
double avg_sum = 0.0;
|
||||||
|
|
||||||
for (int i = 0; i < count; i++) {
|
for (int i = 0; i < count; i++) {
|
||||||
const GpuRecord& r = records[offset + i];
|
const GpuRecord& r = records[offset + i];
|
||||||
|
|
||||||
// min/max
|
// Accumulate avg = (low + high) / 2
|
||||||
if (r.low < stats.low) stats.low = r.low;
|
avg_sum += (r.low + r.high) / 2.0;
|
||||||
if (r.high > stats.high) stats.high = r.high;
|
|
||||||
|
|
||||||
// first/last по timestamp
|
// min/max Open
|
||||||
if (r.timestamp < stats.first_ts) {
|
if (r.open < stats.open_min) stats.open_min = r.open;
|
||||||
stats.first_ts = r.timestamp;
|
if (r.open > stats.open_max) stats.open_max = r.open;
|
||||||
stats.open = r.open;
|
|
||||||
}
|
// min/max Close
|
||||||
if (r.timestamp > stats.last_ts) {
|
if (r.close < stats.close_min) stats.close_min = r.close;
|
||||||
stats.last_ts = r.timestamp;
|
if (r.close > stats.close_max) stats.close_max = r.close;
|
||||||
stats.close = r.close;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
stats.avg = (stats.low + stats.high) / 2.0;
|
stats.avg = avg_sum / static_cast<double>(count);
|
||||||
out_stats[d] = stats;
|
out_stats[d] = stats;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -28,13 +28,17 @@ std::vector<Interval> find_intervals(const std::vector<DayStats>& days, double t
|
|||||||
interval.end_avg = price_now;
|
interval.end_avg = price_now;
|
||||||
interval.change = change;
|
interval.change = change;
|
||||||
|
|
||||||
// Находим min(Open) и max(Close) в интервале
|
// Находим min/max Open и Close в интервале
|
||||||
interval.min_open = days[start_idx].open;
|
interval.open_min = days[start_idx].open_min;
|
||||||
interval.max_close = days[start_idx].close;
|
interval.open_max = days[start_idx].open_max;
|
||||||
|
interval.close_min = days[start_idx].close_min;
|
||||||
|
interval.close_max = days[start_idx].close_max;
|
||||||
|
|
||||||
for (size_t j = start_idx; j <= i; j++) {
|
for (size_t j = start_idx + 1; j <= i; j++) {
|
||||||
interval.min_open = std::min(interval.min_open, days[j].open);
|
interval.open_min = std::min(interval.open_min, days[j].open_min);
|
||||||
interval.max_close = std::max(interval.max_close, days[j].close);
|
interval.open_max = std::max(interval.open_max, days[j].open_max);
|
||||||
|
interval.close_min = std::min(interval.close_min, days[j].close_min);
|
||||||
|
interval.close_max = std::max(interval.close_max, days[j].close_max);
|
||||||
}
|
}
|
||||||
|
|
||||||
intervals.push_back(interval);
|
intervals.push_back(interval);
|
||||||
@@ -68,16 +72,17 @@ void write_intervals(const std::string& filename, const std::vector<Interval>& i
|
|||||||
std::ofstream out(filename);
|
std::ofstream out(filename);
|
||||||
|
|
||||||
out << std::fixed << std::setprecision(2);
|
out << std::fixed << std::setprecision(2);
|
||||||
out << "start_date,end_date,min_open,max_close,start_avg,end_avg,change\n";
|
out << "start_date,end_date,open_min,open_max,close_min,close_max,start_avg,end_avg,change\n";
|
||||||
|
|
||||||
for (const auto& iv : intervals) {
|
for (const auto& iv : intervals) {
|
||||||
out << day_index_to_date(iv.start_day) << ","
|
out << day_index_to_date(iv.start_day) << ","
|
||||||
<< day_index_to_date(iv.end_day) << ","
|
<< day_index_to_date(iv.end_day) << ","
|
||||||
<< iv.min_open << ","
|
<< iv.open_min << ","
|
||||||
<< iv.max_close << ","
|
<< iv.open_max << ","
|
||||||
|
<< iv.close_min << ","
|
||||||
|
<< iv.close_max << ","
|
||||||
<< iv.start_avg << ","
|
<< iv.start_avg << ","
|
||||||
<< iv.end_avg << ","
|
<< iv.end_avg << ","
|
||||||
<< std::setprecision(6) << iv.change << "\n";
|
<< std::setprecision(6) << iv.change << "\n";
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -4,6 +4,19 @@
|
|||||||
#include <vector>
|
#include <vector>
|
||||||
#include <string>
|
#include <string>
|
||||||
|
|
||||||
|
// Интервал с изменением >= threshold
|
||||||
|
struct Interval {
|
||||||
|
DayIndex start_day;
|
||||||
|
DayIndex end_day;
|
||||||
|
double open_min; // минимальный Open в интервале
|
||||||
|
double open_max; // максимальный Open в интервале
|
||||||
|
double close_min; // минимальный Close в интервале
|
||||||
|
double close_max; // максимальный Close в интервале
|
||||||
|
double start_avg;
|
||||||
|
double end_avg;
|
||||||
|
double change;
|
||||||
|
};
|
||||||
|
|
||||||
// Вычисление интервалов с изменением >= threshold (по умолчанию 10%)
|
// Вычисление интервалов с изменением >= threshold (по умолчанию 10%)
|
||||||
std::vector<Interval> find_intervals(const std::vector<DayStats>& days, double threshold = 0.10);
|
std::vector<Interval> find_intervals(const std::vector<DayStats>& days, double threshold = 0.10);
|
||||||
|
|
||||||
@@ -12,4 +25,3 @@ void write_intervals(const std::string& filename, const std::vector<Interval>& i
|
|||||||
|
|
||||||
// Преобразование DayIndex в строку даты (YYYY-MM-DD)
|
// Преобразование DayIndex в строку даты (YYYY-MM-DD)
|
||||||
std::string day_index_to_date(DayIndex day);
|
std::string day_index_to_date(DayIndex day);
|
||||||
|
|
||||||
|
|||||||
22
src/main.cpp
22
src/main.cpp
@@ -4,8 +4,9 @@
|
|||||||
#include <iomanip>
|
#include <iomanip>
|
||||||
|
|
||||||
#include "csv_loader.hpp"
|
#include "csv_loader.hpp"
|
||||||
#include "utils.hpp"
|
|
||||||
#include "record.hpp"
|
#include "record.hpp"
|
||||||
|
#include "day_stats.hpp"
|
||||||
|
#include "aggregation.hpp"
|
||||||
|
|
||||||
int main(int argc, char** argv) {
|
int main(int argc, char** argv) {
|
||||||
MPI_Init(&argc, &argv);
|
MPI_Init(&argc, &argv);
|
||||||
@@ -15,18 +16,27 @@ int main(int argc, char** argv) {
|
|||||||
MPI_Comm_size(MPI_COMM_WORLD, &size);
|
MPI_Comm_size(MPI_COMM_WORLD, &size);
|
||||||
|
|
||||||
// Параллельное чтение данных
|
// Параллельное чтение данных
|
||||||
double start_time = MPI_Wtime();
|
double read_start = MPI_Wtime();
|
||||||
|
|
||||||
std::vector<Record> records = load_csv_parallel(rank, size);
|
std::vector<Record> records = load_csv_parallel(rank, size);
|
||||||
|
double read_time = MPI_Wtime() - read_start;
|
||||||
double end_time = MPI_Wtime();
|
|
||||||
double read_time = end_time - start_time;
|
|
||||||
|
|
||||||
std::cout << "Rank " << rank
|
std::cout << "Rank " << rank
|
||||||
<< ": read " << records.size() << " records"
|
<< ": read " << records.size() << " records"
|
||||||
<< " in " << std::fixed << std::setprecision(3) << read_time << " sec"
|
<< " in " << std::fixed << std::setprecision(3) << read_time << " sec"
|
||||||
<< std::endl;
|
<< std::endl;
|
||||||
|
|
||||||
|
// Агрегация по дням
|
||||||
|
double agg_start = MPI_Wtime();
|
||||||
|
std::vector<DayStats> days = aggregate_days(records);
|
||||||
|
double agg_time = MPI_Wtime() - agg_start;
|
||||||
|
|
||||||
|
std::cout << "Rank " << rank
|
||||||
|
<< ": aggregated " << days.size() << " days"
|
||||||
|
<< " [" << (days.empty() ? 0 : days.front().day)
|
||||||
|
<< ".." << (days.empty() ? 0 : days.back().day) << "]"
|
||||||
|
<< " in " << std::fixed << std::setprecision(3) << agg_time << " sec"
|
||||||
|
<< std::endl;
|
||||||
|
|
||||||
MPI_Finalize();
|
MPI_Finalize();
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user