Line data Source code
1 : /*
2 : * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved.
3 : */
4 : #include "viz_constants.h"
5 : #include "stats_select.h"
6 : #include "stats_query.h"
7 : #include "query.h"
8 : #include <cstdlib>
9 : #include <boost/assign/list_of.hpp>
10 : #include <boost/functional/hash.hpp>
11 : #include "rapidjson/document.h"
12 : #include "rapidjson/stringbuffer.h"
13 : #include "rapidjson/writer.h"
14 :
15 : using boost::assign::map_list_of;
16 :
17 : using std::string;
18 : using std::map;
19 : using std::vector;
20 : using std::set;
21 : using std::pair;
22 : using std::make_pair;
23 :
24 : struct Centroid {
25 : };
26 :
27 : void
28 6 : StatsSelect::DeleteCentroid (Centroid *t) {
29 6 : free(t);
30 6 : }
31 :
32 : struct TDigest {
33 : };
34 :
35 : void
36 3 : StatsSelect::DeleteTDigest(TDigest *t) {
37 3 : TDigest_destroy(t);
38 3 : }
39 :
40 : bool
41 196 : StatsSelect::Jsonify(const std::string& table,
42 : const std::map<std::string, StatVal>& uniks,
43 : const QEOpServerProxy::AggRowT& aggs, std::string& jstr) {
44 :
45 196 : contrail_rapidjson::Document dd;
46 196 : dd.SetObject();
47 :
48 : try {
49 196 : for (std::map<std::string, StatVal>::const_iterator it = uniks.begin();
50 879 : it!= uniks.end(); it++) {
51 683 : switch (it->second.which()) {
52 374 : case QEOpServerProxy::STRING : {
53 374 : string mapit = boost::get<string>(it->second);
54 374 : contrail_rapidjson::Value val(contrail_rapidjson::kStringType);
55 374 : val.SetString(mapit.c_str(), dd.GetAllocator());
56 374 : contrail_rapidjson::Value vk;
57 374 : dd.AddMember(vk.SetString(it->first.c_str(),
58 : dd.GetAllocator()),
59 : val, dd.GetAllocator());
60 374 : }
61 374 : break;
62 49 : case QEOpServerProxy::UUID : {
63 49 : boost::uuids::uuid mapit = boost::get<boost::uuids::uuid>(it->second);
64 49 : contrail_rapidjson::Value val(contrail_rapidjson::kStringType);
65 49 : std::string ustr = to_string(mapit);
66 49 : val.SetString(ustr.c_str(), dd.GetAllocator());
67 49 : contrail_rapidjson::Value vk;
68 49 : dd.AddMember(vk.SetString(it->first.c_str(),
69 : dd.GetAllocator()),
70 : val, dd.GetAllocator());
71 49 : }
72 49 : break;
73 233 : case QEOpServerProxy::UINT64 : {
74 233 : uint64_t mapit = boost::get<uint64_t>(it->second);
75 233 : contrail_rapidjson::Value val(contrail_rapidjson::kNumberType);
76 233 : val.SetUint64(mapit);
77 233 : contrail_rapidjson::Value vk;
78 233 : dd.AddMember(vk.SetString(it->first.c_str(),
79 : dd.GetAllocator()),
80 : val, dd.GetAllocator());
81 233 : }
82 233 : break;
83 27 : case QEOpServerProxy::DOUBLE : {
84 27 : double mapit = boost::get<double>(it->second);
85 27 : contrail_rapidjson::Value val(contrail_rapidjson::kNumberType);
86 27 : val.SetDouble(mapit);
87 27 : contrail_rapidjson::Value vk;
88 27 : dd.AddMember(vk.SetString(it->first.c_str(),
89 : dd.GetAllocator()),
90 : val, dd.GetAllocator());
91 27 : }
92 27 : break;
93 0 : default:
94 0 : QE_ASSERT(0);
95 : }
96 : }
97 196 : for (QEOpServerProxy::AggRowT::const_iterator it = aggs.begin();
98 303 : it!= aggs.end(); it++) {
99 :
100 107 : string sname("");
101 107 : if (it->first.first == QEOpServerProxy::SUM) {
102 159 : if (it->first.second == g_viz_constants.SESSION_FWD_TEARDOWN_BYTES ||
103 150 : it->first.second == g_viz_constants.SESSION_FWD_TEARDOWN_PKTS ||
104 141 : it->first.second == g_viz_constants.SESSION_REV_TEARDOWN_BYTES ||
105 132 : it->first.second == g_viz_constants.SESSION_REV_TEARDOWN_PKTS ||
106 225 : it->first.second == g_viz_constants.FLOW_TABLE_AGG_PKTS ||
107 66 : it->first.second == g_viz_constants.FLOW_TABLE_AGG_BYTES) {
108 18 : sname = it->first.second;
109 : } else {
110 66 : sname = string("SUM(") + it->first.second + string(")");
111 : }
112 23 : } else if (it->first.first == QEOpServerProxy::COUNT) {
113 7 : if (table == g_viz_constants.SESSION_SERIES_TABLE) {
114 5 : sname = "sample_count";
115 : } else {
116 2 : sname = string("COUNT(") + it->first.second + string(")");
117 : }
118 16 : } else if (it->first.first == QEOpServerProxy::CLASS) {
119 9 : if (table == g_viz_constants.SESSION_SERIES_TABLE) {
120 0 : sname = "session_class_id";
121 9 : } else if (table == g_viz_constants.FLOW_SERIES_TABLE){
122 0 : sname = "flow_class_id";
123 : } else {
124 9 : sname = string("CLASS(") + it->first.second + string(")");
125 : }
126 7 : } else if (it->first.first == QEOpServerProxy::MAX) {
127 1 : sname = string("MAX(") + it->first.second + string(")");
128 6 : } else if (it->first.first == QEOpServerProxy::MIN) {
129 1 : sname = string("MIN(") + it->first.second + string(")");
130 5 : } else if (it->first.first == QEOpServerProxy::PERCENTILES) {
131 1 : sname = string("PERCENTILES(") + it->first.second + string(")");
132 4 : } else if (it->first.first == QEOpServerProxy::AVG) {
133 2 : sname = string("AVG(") + it->first.second + string(")");
134 2 : } else if (it->first.first == QEOpServerProxy::COUNT_DISTINCT) {
135 2 : sname = string("COUNT_DISTINCT(") + it->first.second + string(")");
136 : } else {
137 0 : QE_ASSERT(0);
138 : }
139 107 : switch (it->second.which()) {
140 :
141 103 : case QEOpServerProxy::UINT64 : {
142 103 : uint64_t mapit = boost::get<uint64_t>(it->second);
143 103 : contrail_rapidjson::Value val(contrail_rapidjson::kNumberType);
144 103 : val.SetUint64(mapit);
145 103 : contrail_rapidjson::Value vk;
146 103 : dd.AddMember(vk.SetString(sname.c_str(),
147 : dd.GetAllocator()),
148 : val, dd.GetAllocator());
149 103 : }
150 103 : break;
151 1 : case QEOpServerProxy::DOUBLE : {
152 1 : double mapit = boost::get<double>(it->second);
153 1 : contrail_rapidjson::Value val(contrail_rapidjson::kNumberType);
154 1 : val.SetDouble(mapit);
155 1 : contrail_rapidjson::Value vk;
156 1 : dd.AddMember(vk.SetString(sname.c_str(),
157 : dd.GetAllocator()),
158 : val, dd.GetAllocator());
159 1 : }
160 1 : break;
161 1 : case QEOpServerProxy::TDIGEST : {
162 : boost::shared_ptr<TDigest> mapit =
163 1 : boost::get<boost::shared_ptr<TDigest> >(it->second);
164 :
165 1 : contrail_rapidjson::Value val(contrail_rapidjson::kObjectType);
166 1 : float ptiles[7] = {.01,.05,.25,.5,.75,.95,.99};
167 9 : std::string stiles[7] = {"01","05","25","50","75","95","99"};
168 : #if 0
169 : size_t jj = TDigest_get_ncentroids(mapit.get());
170 : for (size_t kk=0; kk < jj; kk++) {
171 : Centroid *c = TDigest_get_centroid(mapit.get(), kk);
172 : std::cout << "Parse Centroid " <<
173 : Centroid_get_mean(c) << " " <<
174 : Centroid_get_count(c) << " " <<
175 : Centroid_quantile(c, mapit.get()) << std::endl;
176 : }
177 : #endif
178 8 : for (size_t idx=0; idx<7; idx++) {
179 7 : contrail_rapidjson::Value sval(contrail_rapidjson::kNumberType);
180 7 : sval.SetDouble(TDigest_percentile(mapit.get(), ptiles[idx]));
181 7 : contrail_rapidjson::Value vk;
182 7 : val.AddMember(vk.SetString(stiles[idx].c_str(),
183 : dd.GetAllocator()),
184 : sval, dd.GetAllocator());
185 7 : }
186 1 : contrail_rapidjson::Value vk;
187 1 : dd.AddMember(vk.SetString(sname.c_str(),
188 : dd.GetAllocator()),
189 : val, dd.GetAllocator());
190 10 : }
191 1 : break;
192 2 : case QEOpServerProxy::CENTROID : {
193 : boost::shared_ptr<Centroid> mapit =
194 2 : boost::get<boost::shared_ptr<Centroid> >(it->second);
195 2 : contrail_rapidjson::Value val(contrail_rapidjson::kNumberType);
196 2 : val.SetDouble(Centroid_get_mean(mapit.get()));
197 2 : contrail_rapidjson::Value vk;
198 2 : dd.AddMember(vk.SetString(sname.c_str(),
199 : dd.GetAllocator()),
200 : val, dd.GetAllocator());
201 2 : }
202 2 : break;
203 0 : default:
204 0 : QE_ASSERT(0);
205 : }
206 107 : }
207 0 : } catch (boost::bad_get& ex) {
208 0 : QE_ASSERT(0);
209 0 : } catch (const std::out_of_range& oor) {
210 0 : QE_ASSERT(0);
211 : }
212 392 : contrail_rapidjson::StringBuffer sb;
213 196 : contrail_rapidjson::Writer<contrail_rapidjson::StringBuffer> writer(sb);
214 196 : dd.Accept(writer);
215 196 : jstr = sb.GetString();
216 : //QE_LOG_NOQID(INFO, "Jsonify row " << jstr);
217 392 : return true;
218 196 : }
219 :
220 64 : void StatsSelect::MergeFinal(const std::vector<boost::shared_ptr<MapBufT> >& inputs,
221 : MapBufT& output) {
222 64 : MapBufT temp;
223 576 : for (size_t idx = 0; idx < inputs.size(); idx++) {
224 512 : if (agg_sort_cols_.empty())
225 504 : Merge(count_distinct_field_, *inputs[idx], output);
226 : else
227 8 : Merge(count_distinct_field_, *inputs[idx], temp);
228 : }
229 64 : if (!agg_sort_cols_.empty()) {
230 1 : MapBufT::iterator jt = temp.end();
231 4 : for (MapBufT::iterator it = temp.begin(); it!= temp.end(); it++) {
232 3 : if (jt!=temp.end()) {
233 2 : temp.erase(jt);
234 : }
235 3 : std::vector<StatVal> ukey = it->first;
236 3 : const StatMap& uniks = it->second.first;
237 3 : const QEOpServerProxy::AggRowT& narows = it->second.second;
238 3 : for (AggSortT::const_iterator xt = agg_sort_cols_.begin();
239 6 : xt != agg_sort_cols_.end(); xt++) {
240 3 : QE_ASSERT(narows.find(xt->first) != narows.end());
241 3 : ukey[xt->second] = narows.at(xt->first);
242 : }
243 3 : MergeFullRow(ukey, uniks, narows, output);
244 3 : jt = it;
245 3 : }
246 1 : temp.clear();
247 : }
248 64 : }
249 :
250 936 : void StatsSelect::Merge(const std::string& count_distinct_field_, const MapBufT& input, MapBufT& output) {
251 936 : for (MapBufT::const_iterator it = input.begin();
252 1490 : it!=input.end(); it++) {
253 554 : std::vector<StatVal> ukey(it->first);
254 554 : StatMap uniks(it->second.first);
255 554 : QEOpServerProxy::AggRowT narows(it->second.second);
256 :
257 554 : StatMap::iterator itr = uniks.find(count_distinct_field_);
258 554 : if (itr != uniks.end()) {
259 : pair<QEOpServerProxy::AggOper, std::string> aggkey(
260 3 : QEOpServerProxy::COUNT_DISTINCT, count_distinct_field_);
261 3 : narows.insert(make_pair(aggkey, (uint64_t) 1));
262 3 : uniks.erase(itr);
263 3 : size_t hash_slot = ukey.size() - 1;
264 3 : uint64_t hash_val = boost::hash_range(uniks.begin(), uniks.end());
265 3 : ukey[hash_slot] = hash_val;
266 3 : }
267 :
268 554 : MergeFullRow(ukey, uniks, narows, output);
269 :
270 554 : }
271 936 : }
272 :
273 1073 : StatsSelect::StatsSelect(AnalyticsQuery * m_query,
274 1073 : const std::vector<std::string> & select_fields) :
275 1073 : main_query(m_query), select_fields_(select_fields),
276 1074 : ts_period_(0), isT_(false), count_field_() {
277 :
278 1073 : status_ = false;
279 :
280 5197 : for (size_t j=0; j<select_fields_.size(); j++) {
281 :
282 4122 : if (select_fields_[j] == g_viz_constants.STAT_TIME_FIELD) {
283 89 : if (ts_period_) {
284 0 : QE_LOG(INFO, "StatsSelect cannot accept both T and T=");
285 0 : return;
286 : }
287 89 : isT_ = true;
288 89 : sort_cols_.insert(std::make_pair(select_fields_[j], 0));
289 89 : QE_TRACE(DEBUG, "StatsSelect T");
290 :
291 4035 : } else if (select_fields_[j].substr(0, g_viz_constants.STAT_TIMEBIN_FIELD.size()) ==
292 : g_viz_constants.STAT_TIMEBIN_FIELD) {
293 40 : if (isT_) {
294 0 : QE_LOG(INFO, "StatsSelect cannot accept both T and T=");
295 0 : return;
296 : }
297 40 : string tsstr = select_fields_[j].substr(g_viz_constants.STAT_TIMEBIN_FIELD.size());
298 40 : ts_period_ = strtoul(tsstr.c_str(), NULL, 10) * 1000000;
299 40 : if (!ts_period_) {
300 0 : QE_LOG(INFO, "StatsSelect cannot accept T= for " << tsstr);
301 0 : return;
302 : }
303 40 : sort_cols_.insert(std::make_pair(g_viz_constants.STAT_TIMEBIN_FIELD,0));
304 40 : QE_TRACE(DEBUG, "StatsSelect T=");
305 40 : } else {
306 3995 : std::string sfield = select_fields_[j];
307 : QEOpServerProxy::AggOper agg =
308 3995 : StatsQuery::ParseAgg(select_fields_[j], sfield);
309 3995 : if (agg == QEOpServerProxy::INVALID) {
310 3474 : unik_cols_.insert(select_fields_[j]);
311 3473 : QE_TRACE(DEBUG, "StatsSelect unik field " << select_fields_[j]);
312 521 : } else if (agg == QEOpServerProxy::COUNT) {
313 54 : count_field_ = sfield;
314 54 : QE_TRACE(DEBUG, "StatsSelect COUNT " << sfield);
315 467 : } else if (agg == QEOpServerProxy::COUNT_DISTINCT) {
316 18 : count_distinct_field_ = sfield;
317 18 : unik_cols_.insert(sfield);
318 18 : QE_TRACE(DEBUG, "StatsSelect COUNT_DISTINCT " << sfield);
319 449 : } else if (agg == QEOpServerProxy::SUM) {
320 288 : sum_cols_.insert(sfield);
321 288 : QE_TRACE(DEBUG, "StatsSelect SUM " << sfield);
322 161 : } else if (agg == QEOpServerProxy::CLASS) {
323 17 : class_cols_.insert(sfield);
324 17 : QE_TRACE(DEBUG, "StatsSelect CLASS " << sfield);
325 144 : } else if (agg == QEOpServerProxy::MAX) {
326 36 : max_field_.insert(sfield);
327 36 : QE_TRACE(DEBUG, "StatsSelect MAX " << sfield);
328 108 : } else if (agg == QEOpServerProxy::MIN) {
329 18 : min_field_.insert(sfield);
330 18 : QE_TRACE(DEBUG, "StatsSelect MIN " << sfield);
331 90 : } else if (agg == QEOpServerProxy::PERCENTILES) {
332 36 : percentile_cols_.insert(sfield);
333 36 : QE_TRACE(DEBUG, "StatsSelect PERCENTILES " << sfield);
334 54 : } else if (agg == QEOpServerProxy::AVG) {
335 54 : avg_field_.insert(sfield);
336 54 : QE_TRACE(DEBUG, "StatsSelect AVG " << sfield);
337 : } else {
338 0 : QE_ASSERT(0);
339 : }
340 3995 : }
341 : }
342 :
343 1075 : status_ = true;
344 0 : }
345 :
346 1075 : void StatsSelect::SetSortOrder(const std::vector<sort_field_t>& sort_fields) {
347 1075 : if (sort_fields.size()) {
348 24 : sort_cols_.clear();
349 : } else {
350 1051 : QE_TRACE(DEBUG, "StatsSelect Sort By T/T=");
351 : }
352 1099 : for (size_t st = 0; st < sort_fields.size(); st++) {
353 :
354 24 : if (sort_fields[st].name.substr(0, g_viz_constants.STAT_TIMEBIN_FIELD.size()) ==
355 : g_viz_constants.STAT_TIMEBIN_FIELD) {
356 0 : QE_TRACE(DEBUG, "StatsSelect Sort By T=");
357 0 : sort_cols_.insert(std::make_pair(g_viz_constants.STAT_TIMEBIN_FIELD,st));
358 24 : } else if (sort_fields[st].name == g_viz_constants.STAT_TIME_FIELD) {
359 0 : QE_TRACE(DEBUG, "StatsSelect Sort By " << sort_fields[st].name);
360 0 : sort_cols_.insert(std::make_pair(sort_fields[st].name, st));
361 : } else {
362 24 : std::string sfield = sort_fields[st].name;
363 : QEOpServerProxy::AggOper agg =
364 24 : StatsQuery::ParseAgg(sort_fields[st].name, sfield);
365 24 : if (main_query->is_stat_table_query(main_query->table()) &&
366 0 : main_query->stats().is_stat_table_static()) {
367 0 : StatsQuery::column_t c = main_query->stats().get_column_desc(sfield);
368 0 : if (c.datatype == QEOpServerProxy::BLANK) {
369 0 : QE_LOG(INFO, "StatsSelect unknown field " << sort_fields[st].name);
370 0 : return;
371 : }
372 0 : }
373 24 : if (agg == QEOpServerProxy::INVALID) {
374 12 : QE_TRACE(DEBUG, "StatsSelect Sort By " << sort_fields[st].name);
375 12 : sort_cols_.insert(std::make_pair(sort_fields[st].name,st));
376 : } else {
377 12 : agg_sort_cols_.insert(make_pair(make_pair(agg, sfield), st));
378 12 : QE_TRACE(DEBUG, "StatsSelect Sort Agg " << agg << " of " << sfield);
379 : }
380 24 : }
381 : }
382 : }
383 :
384 1092 : void StatsSelect::MergeFullRow(
385 : const std::vector<StatVal>& ukey,
386 : const StatMap& uniks,
387 : const QEOpServerProxy::AggRowT& narows,
388 : MapBufT& output) {
389 :
390 1092 : MapBufT::iterator rt = output.find(ukey);
391 1092 : if (rt!= output.end()) {
392 227 : StatMap& rm = rt->second.first;
393 227 : if (rm!=uniks) {
394 0 : output.insert(make_pair(ukey, make_pair(uniks, narows)));
395 : } else {
396 226 : QEOpServerProxy::AggRowT &arows = rt->second.second;
397 226 : MergeAggRow(arows, narows);
398 : }
399 : } else {
400 865 : QEOpServerProxy::AggRowT arows;
401 865 : output.insert(make_pair(ukey, make_pair(uniks, narows)));
402 865 : }
403 1091 : }
404 :
405 :
406 226 : void StatsSelect::MergeAggRow(QEOpServerProxy::AggRowT &arows,
407 : const QEOpServerProxy::AggRowT &narows) {
408 226 : if (arows.empty()) {
409 108 : arows = narows;
410 108 : return;
411 : }
412 118 : for (QEOpServerProxy::AggRowT::iterator jt = arows.begin();
413 361 : jt!= arows.end(); jt++) {
414 243 : QEOpServerProxy::AggRowT::const_iterator kt = narows.find(jt->first);
415 244 : if (kt!=narows.end()) {
416 : // Attribute name must match for aggregate and for the new value
417 208 : QE_ASSERT(jt->first.second == kt->first.second);
418 :
419 207 : if (jt->first.first == QEOpServerProxy::SUM) {
420 173 : StatsSelect::StatVal & sv = jt->second;
421 : try {
422 173 : if (sv.which() == QEOpServerProxy::UINT64) {
423 173 : sv = boost::get<uint64_t>(sv) + boost::get<uint64_t>(kt->second);
424 : }
425 174 : if (sv.which() == QEOpServerProxy::DOUBLE) {
426 0 : sv = boost::get<double>(sv) + boost::get<double>(kt->second);
427 : }
428 0 : } catch (boost::bad_get& ex) {
429 0 : QE_ASSERT(0);
430 0 : } catch (const std::out_of_range& oor) {
431 0 : QE_ASSERT(0);
432 : }
433 :
434 : }
435 :
436 392 : if (jt->first.first == QEOpServerProxy::COUNT ||
437 185 : jt->first.first == QEOpServerProxy::COUNT_DISTINCT) {
438 : try {
439 24 : uint64_t& sv = boost::get<uint64_t>(jt->second);
440 24 : sv = sv + boost::get<uint64_t>(kt->second);
441 0 : } catch (boost::bad_get& ex) {
442 0 : QE_ASSERT(0);
443 0 : } catch (const std::out_of_range& oor) {
444 0 : QE_ASSERT(0);
445 : }
446 : }
447 207 : if (jt->first.first == QEOpServerProxy::MAX) {
448 2 : StatsSelect::StatVal & sv = jt->second;
449 : try {
450 2 : if (sv.which() == QEOpServerProxy::UINT64) {
451 2 : uint64_t& existing_max = boost::get<uint64_t>(jt->second);
452 2 : existing_max = (existing_max > (boost::get<uint64_t>(kt->second))) ? existing_max:(boost::get<uint64_t>(kt->second));
453 : }
454 2 : if (sv.which() == QEOpServerProxy::DOUBLE) {
455 0 : double& existing_max = boost::get<double>(jt->second);
456 0 : existing_max = (existing_max > (boost::get<double>(kt->second))) ? existing_max:(boost::get<double>(kt->second));
457 : }
458 0 : } catch (boost::bad_get& ex) {
459 0 : QE_ASSERT(0);
460 : }
461 : }
462 207 : if (jt->first.first == QEOpServerProxy::MIN) {
463 2 : StatsSelect::StatVal & sv = jt->second;
464 : try {
465 2 : if (sv.which() == QEOpServerProxy::UINT64) {
466 0 : uint64_t& existing_min = boost::get<uint64_t>(jt->second);
467 0 : existing_min = (existing_min < (boost::get<uint64_t>(kt->second))) ? existing_min:(boost::get<uint64_t>(kt->second));
468 : }
469 2 : if (sv.which() == QEOpServerProxy::DOUBLE) {
470 2 : double& existing_min = boost::get<double>(jt->second);
471 2 : existing_min = (existing_min < (boost::get<double>(kt->second))) ? existing_min:(boost::get<double>(kt->second));
472 : }
473 0 : } catch (boost::bad_get& ex) {
474 0 : QE_ASSERT(0);
475 : }
476 : }
477 207 : if (jt->first.first == QEOpServerProxy::AVG) {
478 4 : StatsSelect::StatVal & sv = jt->second;
479 : try {
480 4 : if (sv.which() == QEOpServerProxy::CENTROID) {
481 : boost::shared_ptr<Centroid>& pt =
482 4 : boost::get<boost::shared_ptr<Centroid> >(jt->second);
483 :
484 : boost::shared_ptr<Centroid> pt2 =
485 4 : boost::get<boost::shared_ptr<Centroid> >(kt->second);
486 4 : Centroid_add(pt.get(),
487 : Centroid_get_mean(pt2.get()),
488 : Centroid_get_count(pt2.get()));
489 4 : }
490 0 : } catch (boost::bad_get& ex) {
491 0 : QE_ASSERT(0);
492 : }
493 : }
494 207 : if (jt->first.first == QEOpServerProxy::PERCENTILES) {
495 2 : StatsSelect::StatVal & sv = jt->second;
496 : try {
497 2 : if (sv.which() == QEOpServerProxy::TDIGEST) {
498 : boost::shared_ptr<TDigest>& pt =
499 2 : boost::get<boost::shared_ptr<TDigest> >(jt->second);
500 :
501 : boost::shared_ptr<TDigest> pt2 =
502 2 : boost::get<boost::shared_ptr<TDigest> >(kt->second);
503 2 : size_t j, ncentroids = TDigest_get_ncentroids(pt2.get());
504 4 : for (j = 0; j < ncentroids; j++) {
505 2 : Centroid * c = TDigest_get_centroid(pt2.get(), j);
506 2 : TDigest* nd = TDigest_add(pt.get(),
507 : Centroid_get_mean(c),
508 : Centroid_get_count(c));
509 2 : if (nd) {
510 0 : pt.reset(nd);
511 : }
512 : }
513 2 : }
514 0 : } catch (boost::bad_get& ex) {
515 0 : QE_ASSERT(0);
516 : }
517 : }
518 : }
519 : }
520 : }
521 :
522 : namespace boost {
523 : std::size_t hash_value(const StatsSelect::StatVal&);
524 : }
525 :
526 1822 : std::size_t boost::hash_value(const StatsSelect::StatVal& sv) {
527 1822 : std::ostringstream ostr;
528 1822 : ostr << sv;
529 3643 : return boost::hash_value(ostr.str());
530 1821 : }
531 :
532 535 : bool StatsSelect::LoadRow(boost::uuids::uuid u,
533 : uint64_t timestamp, const vector<StatEntry>& row, MapBufT& output) {
534 :
535 535 : if (!Status()) return false;
536 535 : uint64_t ts = 0;
537 :
538 : // Build Uniks map
539 535 : StatMap uniks;
540 : set<string>::const_iterator ukit =
541 535 : unik_cols_.find(g_viz_constants.STAT_UUID_FIELD);
542 535 : if (ukit!=unik_cols_.end()) {
543 63 : uniks.insert(make_pair(g_viz_constants.STAT_UUID_FIELD,u));
544 : }
545 :
546 535 : if (isT_) {
547 50 : ts = timestamp;
548 50 : uniks.insert(make_pair(g_viz_constants.STAT_TIME_FIELD,ts));
549 : }
550 :
551 535 : if (ts_period_) {
552 54 : ts = timestamp - (timestamp % ts_period_);
553 54 : uniks.insert(make_pair(g_viz_constants.STAT_TIMEBIN_FIELD,ts));
554 : }
555 :
556 535 : for (vector<StatEntry>::const_iterator it = row.begin();
557 18104 : it != row.end(); it++) {
558 17571 : set<string>::const_iterator uit = unik_cols_.find(it->name);
559 17568 : if (uit!=unik_cols_.end()) {
560 1633 : uniks.insert(make_pair(it->name, it->value));
561 : }
562 : }
563 :
564 : // Build sort vector
565 : // Last slot is reserved for the hash
566 535 : std::vector<StatVal> ukey(sort_cols_.size() + agg_sort_cols_.size() + 1);
567 534 : size_t hash_slot = sort_cols_.size() + agg_sort_cols_.size();
568 534 : uint64_t hash_val = boost::hash_range(uniks.begin(), uniks.end());
569 535 : ukey[hash_slot] = hash_val;
570 :
571 534 : for (map<string, size_t>::const_iterator st = sort_cols_.begin();
572 664 : st!=sort_cols_.end(); st++) {
573 130 : QE_ASSERT(uniks.find(st->first) != uniks.end());
574 130 : ukey[st->second] = uniks.at(st->first);
575 : }
576 :
577 534 : QEOpServerProxy::AggRowT narows;
578 534 : for (vector<StatEntry>::const_iterator it = row.begin();
579 18115 : it != row.end(); it++) {
580 17577 : set<string>::const_iterator uit = sum_cols_.find(it->name);
581 17580 : if (uit!=sum_cols_.end()) {
582 270 : pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::SUM,it->name);
583 270 : narows.insert(make_pair(aggkey, it->value));
584 271 : }
585 : }
586 :
587 535 : for (vector<StatEntry>::const_iterator it = row.begin();
588 18129 : it != row.end(); it++) {
589 17594 : set<string>::const_iterator uit = max_field_.find(it->name);
590 17593 : if (uit!=max_field_.end()) {
591 3 : pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::MAX,it->name);
592 3 : narows.insert(make_pair(aggkey, it->value));
593 3 : }
594 17590 : uit = min_field_.find(it->name);
595 17589 : if (uit!=min_field_.end()) {
596 3 : pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::MIN,it->name);
597 3 : narows.insert(make_pair(aggkey, it->value));
598 3 : }
599 17589 : uit = avg_field_.find(it->name);
600 17595 : if (uit!=avg_field_.end()) {
601 6 : pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::AVG,it->name);
602 : try {
603 6 : Centroid *t = Centroid_create(0, 0);
604 6 : StatsSelect::StatVal sv = it->value;
605 : double val;
606 6 : if (sv.which() == QEOpServerProxy::UINT64) {
607 3 : val = boost::get<uint64_t>(sv);
608 : } else {
609 3 : val = boost::get<double>(sv);
610 : }
611 6 : Centroid_add(t, val, 1);
612 6 : boost::shared_ptr<Centroid> pt(t, &StatsSelect::DeleteCentroid);
613 6 : narows.insert(make_pair(aggkey, pt));
614 6 : } catch (boost::bad_get& ex) {
615 0 : QE_ASSERT(0);
616 0 : } catch (const std::out_of_range& oor) {
617 0 : QE_ASSERT(0);
618 : }
619 6 : }
620 17594 : uit = percentile_cols_.find(it->name);
621 17594 : if (uit!=percentile_cols_.end()) {
622 3 : pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::PERCENTILES,it->name);
623 :
624 : try {
625 3 : TDigest *t = TDigest_create(0.01, 100);
626 3 : StatsSelect::StatVal sv = it->value;
627 : double val;
628 3 : if (sv.which() == QEOpServerProxy::UINT64) {
629 3 : val = boost::get<uint64_t>(sv);
630 : } else {
631 0 : val = boost::get<double>(sv);
632 : }
633 3 : TDigest_add(t, val, 1);
634 3 : boost::shared_ptr<TDigest> pt(t, &StatsSelect::DeleteTDigest);
635 3 : narows.insert(make_pair(aggkey, pt));
636 3 : } catch (boost::bad_get& ex) {
637 0 : QE_ASSERT(0);
638 0 : } catch (const std::out_of_range& oor) {
639 0 : QE_ASSERT(0);
640 : }
641 3 : }
642 : }
643 :
644 535 : for (std::set<std::string>::const_iterator ct = class_cols_.begin();
645 544 : ct!=class_cols_.end(); ct++) {
646 9 : pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::CLASS,*ct);
647 9 : StatMap huniks;
648 9 : for (vector<StatEntry>::const_iterator rit = row.begin();
649 90 : rit != row.end(); rit++) {
650 81 : if (rit->name != *ct) {
651 81 : if (uniks.find(rit->name) != uniks.end()) {
652 : // For generating the hash, consider all attributes that
653 : // are in the row, and that do not match the CLASS attribute,
654 : // and that are in non-aggregate attributes in the SELECT
655 18 : huniks[rit->name] = rit->value;
656 : }
657 : }
658 : }
659 9 : uint64_t hh = boost::hash_range(huniks.begin(), huniks.end());
660 9 : narows.insert(make_pair(aggkey, hh));
661 9 : }
662 :
663 535 : if (!count_field_.empty()) {
664 30 : pair<QEOpServerProxy::AggOper,string> aggkey(QEOpServerProxy::COUNT,count_field_);
665 30 : narows.insert(make_pair(aggkey, (uint64_t) 1));
666 30 : }
667 :
668 535 : MergeFullRow(ukey, uniks, narows, output);
669 :
670 534 : return true;
671 534 : }
672 :
|