-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmain.cpp
More file actions
120 lines (102 loc) · 4.54 KB
/
Copy pathmain.cpp
File metadata and controls
120 lines (102 loc) · 4.54 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
#include <iostream>
#include <string>
#include <vector>
#include <random>
#include <thread>
#include <chrono>
#include <filesystem>
#include <atomic>
#include "events.h"
#include "kafka_sim.h"
#include "clickhouse_sim.h"
#include "metrics_agg.h"
#include "prometheus_server.h"
using namespace std::chrono_literals;
struct Config {
std::string data_dir = "data";
std::string kafka_log = "data/kafka/topic.tsv";
std::string checkpoint = "data/checkpoint.txt";
std::string ch_base = "data/ch";
std::string prom_addr = "127.0.0.1";
unsigned short prom_port = 9090;
int duration_seconds = 15; // demo by default
int eps = 200; // events per second
uint32_t seed = 42;
};
static Config parse_args(int argc, char** argv) {
Config c;
for (int i = 1; i < argc; ++i) {
std::string a = argv[i];
auto next = [&](int i){ return i + 1 < argc ? std::string(argv[i+1]) : std::string(); };
if (a == "--duration-seconds" && i + 1 < argc) { c.duration_seconds = std::stoi(next(i)); ++i; }
else if (a == "--prom-port" && i + 1 < argc) { c.prom_port = static_cast<unsigned short>(std::stoi(next(i))); ++i; }
else if (a == "--prom-addr" && i + 1 < argc) { c.prom_addr = next(i); ++i; }
else if (a == "--eps" && i + 1 < argc) { c.eps = std::stoi(next(i)); ++i; }
else if (a == "--seed" && i + 1 < argc) { c.seed = static_cast<uint32_t>(std::stoul(next(i))); ++i; }
else if (a == "--data-dir" && i + 1 < argc) { c.data_dir = next(i); ++i; }
}
// derive paths if data_dir changed
c.kafka_log = c.data_dir + "/kafka/topic.tsv";
c.checkpoint = c.data_dir + "/checkpoint.txt";
c.ch_base = c.data_dir + "/ch";
return c;
}
static void ensure_dirs(const Config& c) {
namespace fs = std::filesystem;
fs::create_directories(fs::path(c.kafka_log).parent_path());
fs::create_directories(fs::path(c.ch_base));
fs::create_directories(fs::path(c.checkpoint).parent_path());
}
static void producer_thread(const Config& c, adp::KafkaTopic& topic, std::atomic<bool>& stop_flag) {
std::mt19937 rng(c.seed);
std::vector<std::string> campaigns = {"cmp-1", "cmp-2", "cmp-3"};
std::vector<std::string> ads = {"ad-1", "ad-2", "ad-3"};
std::uniform_int_distribution<int> cmpd(0, (int)campaigns.size()-1);
std::uniform_int_distribution<int> add(0, (int)ads.size()-1);
std::uniform_real_distribution<double> u(0.0, 1.0);
const auto period = std::chrono::microseconds(1'000'000 / std::max(1, c.eps));
auto next_tick = std::chrono::steady_clock::now();
while (!stop_flag.load()) {
next_tick += period;
adp::AdEvent e;
e.ts_ms = (uint64_t)std::chrono::duration_cast<std::chrono::milliseconds>(std::chrono::system_clock::now().time_since_epoch()).count();
e.campaign_id = campaigns[cmpd(rng)];
e.ad_id = ads[add(rng)];
double r = u(rng);
if (r < 0.80) { // 80% impressions
e.type = adp::EventType::Impression; e.cost = 0.001; e.revenue = 0.0;
} else if (r < 0.95) { // 15% clicks
e.type = adp::EventType::Click; e.cost = 0.050; e.revenue = 0.0;
} else { // 5% conversions
e.type = adp::EventType::Conversion; e.cost = 0.0; e.revenue = 1.500;
}
topic.append(e);
std::this_thread::sleep_until(next_tick);
}
}
int main(int argc, char** argv) {
Config cfg = parse_args(argc, argv);
ensure_dirs(cfg);
std::cout << "Starting ad-pipeline demo for " << cfg.duration_seconds << "s, EPS=" << cfg.eps << "\n";
adp::KafkaTopic topic(cfg.kafka_log);
adp::ClickHouseSim ch(cfg.ch_base);
adp::MetricsAggregator agg(topic, ch, cfg.checkpoint);
// Prometheus server exposing aggregator metrics
adp::PrometheusServer prom(cfg.prom_addr, cfg.prom_port, [&agg]{ return agg.export_metrics_text(); });
// Start components
agg.start();
prom.start();
std::atomic<bool> stop_prod{false};
std::thread prod([&]{ producer_thread(cfg, topic, stop_prod); });
// Let it run for duration
std::this_thread::sleep_for(std::chrono::seconds(cfg.duration_seconds));
// Shutdown
stop_prod.store(true);
if (prod.joinable()) prod.join();
agg.stop();
prom.stop();
std::cout << "Processed events: " << agg.processed() << "\n";
std::cout << "Prometheus on http://" << cfg.prom_addr << ":" << cfg.prom_port << "/metrics (if running longer)." << std::endl;
std::cout << "Raw events: " << cfg.ch_base << "/raw/events.tsv; Aggregates: " << cfg.ch_base << "/agg/snapshot.tsv" << std::endl;
return 0;
}