From 34ea1103a0faa0e01df0e2b0ac1ce98e7ec3e3f1 Mon Sep 17 00:00:00 2001 From: Yangster-mac Date: Mon, 29 Jan 2018 12:40:55 -0800 Subject: [PATCH] Extend gauge metric to support memory metric. Test: statd unit test passed. Test: statsd unit test passed Change-Id: I2e3f26563678ae77d44afe168454b6d1ea449f3a --- .../src/metrics/GaugeMetricProducer.cpp | 76 +++++++++++++------ cmds/statsd/src/metrics/GaugeMetricProducer.h | 17 +++-- cmds/statsd/src/stats_log.proto | 4 +- cmds/statsd/src/statsd_config.proto | 6 ++ .../statsd/tests/e2e/GaugeMetric_e2e_test.cpp | 28 ++++--- .../metrics/GaugeMetricProducer_test.cpp | 29 ++++--- 6 files changed, 108 insertions(+), 52 deletions(-) diff --git a/cmds/statsd/src/metrics/GaugeMetricProducer.cpp b/cmds/statsd/src/metrics/GaugeMetricProducer.cpp index 24dc5b0fba532..1072c5aae6e40 100644 --- a/cmds/statsd/src/metrics/GaugeMetricProducer.cpp +++ b/cmds/statsd/src/metrics/GaugeMetricProducer.cpp @@ -58,6 +58,7 @@ const int FIELD_ID_BUCKET_INFO = 3; const int FIELD_ID_START_BUCKET_NANOS = 1; const int FIELD_ID_END_BUCKET_NANOS = 2; const int FIELD_ID_ATOM = 3; +const int FIELD_ID_TIMESTAMP = 4; GaugeMetricProducer::GaugeMetricProducer(const ConfigKey& key, const GaugeMetric& metric, const int conditionIndex, @@ -67,7 +68,7 @@ GaugeMetricProducer::GaugeMetricProducer(const ConfigKey& key, const GaugeMetric : MetricProducer(metric.id(), key, startTimeNs, conditionIndex, wizard), mStatsPullerManager(statsPullerManager), mPullTagId(pullTagId) { - mCurrentSlicedBucket = std::make_shared(); + mCurrentSlicedBucket = std::make_shared(); mCurrentSlicedBucketForAnomaly = std::make_shared(); int64_t bucketSizeMills = 0; if (metric.has_bucket()) { @@ -77,6 +78,7 @@ GaugeMetricProducer::GaugeMetricProducer(const ConfigKey& key, const GaugeMetric } mBucketSizeNs = bucketSizeMills * 1000000; + mSamplingType = metric.sampling_type(); mFieldFilter = metric.gauge_fields_filter(); // TODO: use UidMap if uid->pkg_name is required @@ -89,7 +91,7 @@ GaugeMetricProducer::GaugeMetricProducer(const ConfigKey& key, const GaugeMetric } // Kicks off the puller immediately. - if (mPullTagId != -1) { + if (mPullTagId != -1 && mSamplingType == GaugeMetric::RANDOM_ONE_SAMPLE) { mStatsPullerManager->RegisterReceiver(mPullTagId, this, bucketSizeMills); } @@ -154,12 +156,23 @@ void GaugeMetricProducer::onDumpReportLocked(const uint64_t dumpTimeNs, (long long)bucket.mBucketStartNs); protoOutput->write(FIELD_TYPE_INT64 | FIELD_ID_END_BUCKET_NANOS, (long long)bucket.mBucketEndNs); - long long atomToken = protoOutput->start(FIELD_TYPE_MESSAGE | FIELD_ID_ATOM); - writeFieldValueTreeToStream(*bucket.mGaugeFields, protoOutput); - protoOutput->end(atomToken); + + if (!bucket.mGaugeAtoms.empty()) { + long long atomsToken = + protoOutput->start(FIELD_TYPE_MESSAGE | FIELD_COUNT_REPEATED | FIELD_ID_ATOM); + for (const auto& atom : bucket.mGaugeAtoms) { + writeFieldValueTreeToStream(*atom.mFields, protoOutput); + } + protoOutput->end(atomsToken); + + for (const auto& atom : bucket.mGaugeAtoms) { + protoOutput->write(FIELD_TYPE_INT64 | FIELD_COUNT_REPEATED | FIELD_ID_TIMESTAMP, + (long long)atom.mTimestamps); + } + } protoOutput->end(bucketInfoToken); - VLOG("\t bucket [%lld - %lld] includes %d gauge fields.", (long long)bucket.mBucketStartNs, - (long long)bucket.mBucketEndNs, (int)bucket.mGaugeFields->size()); + VLOG("\t bucket [%lld - %lld] includes %d atoms.", (long long)bucket.mBucketStartNs, + (long long)bucket.mBucketEndNs, (int)bucket.mGaugeAtoms.size()); } protoOutput->end(wrapperToken); } @@ -181,14 +194,26 @@ void GaugeMetricProducer::onConditionChangedLocked(const bool conditionMet, if (mPullTagId == -1) { return; } - // No need to pull again. Either scheduled pull or condition on true happened - if (!mCondition) { - return; - } - // Already have gauge metric for the current bucket, do not do it again. - if (mCurrentSlicedBucket->size() > 0) { + + bool triggerPuller = false; + switch(mSamplingType) { + // When the metric wants to do random sampling and there is already one gauge atom for the + // current bucket, do not do it again. + case GaugeMetric::RANDOM_ONE_SAMPLE: { + triggerPuller = mCondition && mCurrentSlicedBucket->empty(); + break; + } + case GaugeMetric::ALL_CONDITION_CHANGES: { + triggerPuller = true; + break; + } + default: + break; + } + if (!triggerPuller) { return; } + vector> allData; if (!mStatsPullerManager->Pull(mPullTagId, &allData)) { ALOGE("Stats puller failed for tag: %d", mPullTagId); @@ -257,20 +282,24 @@ void GaugeMetricProducer::onMatchedLogEventInternalLocked( } flushIfNeededLocked(eventTimeNs); - // For gauge metric, we just simply use the first gauge in the given bucket. - if (mCurrentSlicedBucket->find(eventKey) != mCurrentSlicedBucket->end()) { + // When gauge metric wants to randomly sample the output atom, we just simply use the first + // gauge in the given bucket. + if (mCurrentSlicedBucket->find(eventKey) != mCurrentSlicedBucket->end() && + mSamplingType == GaugeMetric::RANDOM_ONE_SAMPLE) { return; } - std::shared_ptr gaugeFields = getGaugeFields(event); if (hitGuardRailLocked(eventKey)) { return; } - (*mCurrentSlicedBucket)[eventKey] = gaugeFields; + GaugeAtom gaugeAtom; + gaugeAtom.mFields = getGaugeFields(event); + gaugeAtom.mTimestamps = eventTimeNs; + (*mCurrentSlicedBucket)[eventKey].push_back(gaugeAtom); // Anomaly detection on gauge metric only works when there is one numeric // field specified. if (mAnomalyTrackers.size() > 0) { - if (gaugeFields->size() == 1) { - const DimensionsValue& dimensionsValue = gaugeFields->begin()->second; + if (gaugeAtom.mFields->size() == 1) { + const DimensionsValue& dimensionsValue = gaugeAtom.mFields->begin()->second; long gaugeVal = 0; if (dimensionsValue.has_value_int()) { gaugeVal = (long)dimensionsValue.value_int(); @@ -289,7 +318,10 @@ void GaugeMetricProducer::updateCurrentSlicedBucketForAnomaly() { mCurrentSlicedBucketForAnomaly->clear(); status_t err = NO_ERROR; for (const auto& slice : *mCurrentSlicedBucket) { - const DimensionsValue& dimensionsValue = slice.second->begin()->second; + if (slice.second.empty() || slice.second.front().mFields->empty()) { + continue; + } + const DimensionsValue& dimensionsValue = slice.second.front().mFields->begin()->second; long gaugeVal = 0; if (dimensionsValue.has_value_int()) { gaugeVal = (long)dimensionsValue.value_int(); @@ -318,7 +350,7 @@ void GaugeMetricProducer::flushIfNeededLocked(const uint64_t& eventTimeNs) { info.mBucketNum = mCurrentBucketNum; for (const auto& slice : *mCurrentSlicedBucket) { - info.mGaugeFields = slice.second; + info.mGaugeAtoms = slice.second; auto& bucketList = mPastBuckets[slice.first]; bucketList.push_back(info); VLOG("gauge metric %lld, dump key value: %s", @@ -334,7 +366,7 @@ void GaugeMetricProducer::flushIfNeededLocked(const uint64_t& eventTimeNs) { } mCurrentSlicedBucketForAnomaly = std::make_shared(); - mCurrentSlicedBucket = std::make_shared(); + mCurrentSlicedBucket = std::make_shared(); // Adjusts the bucket start time int64_t numBucketsForward = (eventTimeNs - mCurrentBucketStartTimeNs) / mBucketSizeNs; diff --git a/cmds/statsd/src/metrics/GaugeMetricProducer.h b/cmds/statsd/src/metrics/GaugeMetricProducer.h index 1895edf8f7291..6c013477af374 100644 --- a/cmds/statsd/src/metrics/GaugeMetricProducer.h +++ b/cmds/statsd/src/metrics/GaugeMetricProducer.h @@ -32,15 +32,20 @@ namespace android { namespace os { namespace statsd { +struct GaugeAtom { + std::shared_ptr mFields; + int64_t mTimestamps; +}; + struct GaugeBucket { int64_t mBucketStartNs; int64_t mBucketEndNs; - std::shared_ptr mGaugeFields; + std::vector mGaugeAtoms; uint64_t mBucketNum; }; -typedef std::unordered_map> - DimToGaugeFieldsMap; +typedef std::unordered_map> + DimToGaugeAtomsMap; // This gauge metric producer first register the puller to automatically pull the gauge at the // beginning of each bucket. If the condition is met, insert it to the bucket info. Otherwise @@ -48,7 +53,7 @@ typedef std::unordered_map> // producer always reports the guage at the earliest time of the bucket when the condition is met. class GaugeMetricProducer : public virtual MetricProducer, public virtual PullDataReceiver { public: - GaugeMetricProducer(const ConfigKey& key, const GaugeMetric& countMetric, + GaugeMetricProducer(const ConfigKey& key, const GaugeMetric& gaugeMetric, const int conditionIndex, const sp& wizard, const int pullTagId, const int64_t startTimeNs); @@ -97,7 +102,7 @@ private: std::unordered_map> mPastBuckets; // The current bucket. - std::shared_ptr mCurrentSlicedBucket; + std::shared_ptr mCurrentSlicedBucket; // The current bucket for anomaly detection. std::shared_ptr mCurrentSlicedBucketForAnomaly; @@ -108,6 +113,8 @@ private: // Whitelist of fields to report. Empty means all are reported. FieldFilter mFieldFilter; + GaugeMetric::SamplingType mSamplingType; + // apply a whitelist on the original input std::shared_ptr getGaugeFields(const LogEvent& event); diff --git a/cmds/statsd/src/stats_log.proto b/cmds/statsd/src/stats_log.proto index a4ccbd42883da..af21ca4d82e87 100644 --- a/cmds/statsd/src/stats_log.proto +++ b/cmds/statsd/src/stats_log.proto @@ -100,7 +100,9 @@ message GaugeBucketInfo { optional int64 end_bucket_nanos = 2; - optional Atom atom = 3; + repeated Atom atom = 3; + + repeated int64 timestamp_nanos = 4; } message GaugeMetricData { diff --git a/cmds/statsd/src/statsd_config.proto b/cmds/statsd/src/statsd_config.proto index 07bbcb2190e8d..2ea79a64a5ea4 100644 --- a/cmds/statsd/src/statsd_config.proto +++ b/cmds/statsd/src/statsd_config.proto @@ -222,6 +222,12 @@ message GaugeMetric { optional TimeUnit bucket = 6; repeated MetricConditionLink links = 7; + + enum SamplingType { + RANDOM_ONE_SAMPLE = 1; + ALL_CONDITION_CHANGES = 2; + } + optional SamplingType sampling_type = 9 [default = RANDOM_ONE_SAMPLE] ; } message ValueMetric { diff --git a/cmds/statsd/tests/e2e/GaugeMetric_e2e_test.cpp b/cmds/statsd/tests/e2e/GaugeMetric_e2e_test.cpp index e56a6c57848f2..a80fdc5606b75 100644 --- a/cmds/statsd/tests/e2e/GaugeMetric_e2e_test.cpp +++ b/cmds/statsd/tests/e2e/GaugeMetric_e2e_test.cpp @@ -153,23 +153,26 @@ TEST(GaugeMetricE2eTest, TestMultipleFieldsForPushedEvent) { EXPECT_EQ(data.dimensions_in_what().value_tuple().dimensions_value(0).field(), 1 /* uid field */); EXPECT_EQ(data.dimensions_in_what().value_tuple().dimensions_value(0).value_int(), appUid1); EXPECT_EQ(data.bucket_info_size(), 3); + EXPECT_EQ(data.bucket_info(0).atom_size(), 1); EXPECT_EQ(data.bucket_info(0).start_bucket_nanos(), bucketStartTimeNs); EXPECT_EQ(data.bucket_info(0).end_bucket_nanos(), bucketStartTimeNs + bucketSizeNs); - EXPECT_EQ(data.bucket_info(0).atom().app_start_changed().type(), AppStartChanged::HOT); - EXPECT_EQ(data.bucket_info(0).atom().app_start_changed().activity_name(), "activity_name2"); - EXPECT_EQ(data.bucket_info(0).atom().app_start_changed().activity_start_msec(), 102L); + EXPECT_EQ(data.bucket_info(0).atom(0).app_start_changed().type(), AppStartChanged::HOT); + EXPECT_EQ(data.bucket_info(0).atom(0).app_start_changed().activity_name(), "activity_name2"); + EXPECT_EQ(data.bucket_info(0).atom(0).app_start_changed().activity_start_msec(), 102L); + EXPECT_EQ(data.bucket_info(1).atom_size(), 1); EXPECT_EQ(data.bucket_info(1).start_bucket_nanos(), bucketStartTimeNs + bucketSizeNs); EXPECT_EQ(data.bucket_info(1).end_bucket_nanos(), bucketStartTimeNs + 2 * bucketSizeNs); - EXPECT_EQ(data.bucket_info(1).atom().app_start_changed().type(), AppStartChanged::WARM); - EXPECT_EQ(data.bucket_info(1).atom().app_start_changed().activity_name(), "activity_name4"); - EXPECT_EQ(data.bucket_info(1).atom().app_start_changed().activity_start_msec(), 104L); + EXPECT_EQ(data.bucket_info(1).atom(0).app_start_changed().type(), AppStartChanged::WARM); + EXPECT_EQ(data.bucket_info(1).atom(0).app_start_changed().activity_name(), "activity_name4"); + EXPECT_EQ(data.bucket_info(1).atom(0).app_start_changed().activity_start_msec(), 104L); + EXPECT_EQ(data.bucket_info(2).atom_size(), 1); EXPECT_EQ(data.bucket_info(2).start_bucket_nanos(), bucketStartTimeNs + 2 * bucketSizeNs); EXPECT_EQ(data.bucket_info(2).end_bucket_nanos(), bucketStartTimeNs + 3 * bucketSizeNs); - EXPECT_EQ(data.bucket_info(2).atom().app_start_changed().type(), AppStartChanged::COLD); - EXPECT_EQ(data.bucket_info(2).atom().app_start_changed().activity_name(), "activity_name5"); - EXPECT_EQ(data.bucket_info(2).atom().app_start_changed().activity_start_msec(), 105L); + EXPECT_EQ(data.bucket_info(2).atom(0).app_start_changed().type(), AppStartChanged::COLD); + EXPECT_EQ(data.bucket_info(2).atom(0).app_start_changed().activity_name(), "activity_name5"); + EXPECT_EQ(data.bucket_info(2).atom(0).app_start_changed().activity_start_msec(), 105L); data = gaugeMetrics.data(1); @@ -178,11 +181,12 @@ TEST(GaugeMetricE2eTest, TestMultipleFieldsForPushedEvent) { EXPECT_EQ(data.dimensions_in_what().value_tuple().dimensions_value(0).field(), 1 /* uid field */); EXPECT_EQ(data.dimensions_in_what().value_tuple().dimensions_value(0).value_int(), appUid2); EXPECT_EQ(data.bucket_info_size(), 1); + EXPECT_EQ(data.bucket_info(0).atom_size(), 1); EXPECT_EQ(data.bucket_info(0).start_bucket_nanos(), bucketStartTimeNs + 2 * bucketSizeNs); EXPECT_EQ(data.bucket_info(0).end_bucket_nanos(), bucketStartTimeNs + 3 * bucketSizeNs); - EXPECT_EQ(data.bucket_info(0).atom().app_start_changed().type(), AppStartChanged::COLD); - EXPECT_EQ(data.bucket_info(0).atom().app_start_changed().activity_name(), "activity_name7"); - EXPECT_EQ(data.bucket_info(0).atom().app_start_changed().activity_start_msec(), 201L); + EXPECT_EQ(data.bucket_info(0).atom(0).app_start_changed().type(), AppStartChanged::COLD); + EXPECT_EQ(data.bucket_info(0).atom(0).app_start_changed().activity_name(), "activity_name7"); + EXPECT_EQ(data.bucket_info(0).atom(0).app_start_changed().activity_start_msec(), 201L); } #else diff --git a/cmds/statsd/tests/metrics/GaugeMetricProducer_test.cpp b/cmds/statsd/tests/metrics/GaugeMetricProducer_test.cpp index 82772d854db27..4533ac610057b 100644 --- a/cmds/statsd/tests/metrics/GaugeMetricProducer_test.cpp +++ b/cmds/statsd/tests/metrics/GaugeMetricProducer_test.cpp @@ -78,7 +78,7 @@ TEST(GaugeMetricProducerTest, TestNoCondition) { gaugeProducer.onDataPulled(allData); EXPECT_EQ(1UL, gaugeProducer.mCurrentSlicedBucket->size()); - auto it = gaugeProducer.mCurrentSlicedBucket->begin()->second->begin(); + auto it = gaugeProducer.mCurrentSlicedBucket->begin()->second.front().mFields->begin(); EXPECT_EQ(10, it->second.value_int()); it++; EXPECT_EQ(11, it->second.value_int()); @@ -94,14 +94,14 @@ TEST(GaugeMetricProducerTest, TestNoCondition) { allData.push_back(event2); gaugeProducer.onDataPulled(allData); EXPECT_EQ(1UL, gaugeProducer.mCurrentSlicedBucket->size()); - it = gaugeProducer.mCurrentSlicedBucket->begin()->second->begin(); + it = gaugeProducer.mCurrentSlicedBucket->begin()->second.front().mFields->begin(); EXPECT_EQ(24, it->second.value_int()); it++; EXPECT_EQ(25, it->second.value_int()); // One dimension. EXPECT_EQ(1UL, gaugeProducer.mPastBuckets.size()); EXPECT_EQ(1UL, gaugeProducer.mPastBuckets.begin()->second.size()); - it = gaugeProducer.mPastBuckets.begin()->second.back().mGaugeFields->begin(); + it = gaugeProducer.mPastBuckets.begin()->second.back().mGaugeAtoms.front().mFields->begin(); EXPECT_EQ(10L, it->second.value_int()); it++; EXPECT_EQ(11L, it->second.value_int()); @@ -112,7 +112,7 @@ TEST(GaugeMetricProducerTest, TestNoCondition) { // One dimension. EXPECT_EQ(1UL, gaugeProducer.mPastBuckets.size()); EXPECT_EQ(2UL, gaugeProducer.mPastBuckets.begin()->second.size()); - it = gaugeProducer.mPastBuckets.begin()->second.back().mGaugeFields->begin(); + it = gaugeProducer.mPastBuckets.begin()->second.back().mGaugeAtoms.front().mFields->begin(); EXPECT_EQ(24L, it->second.value_int()); it++; EXPECT_EQ(25L, it->second.value_int()); @@ -151,7 +151,8 @@ TEST(GaugeMetricProducerTest, TestWithCondition) { gaugeProducer.onConditionChanged(true, bucketStartTimeNs + 8); EXPECT_EQ(1UL, gaugeProducer.mCurrentSlicedBucket->size()); EXPECT_EQ(100, - gaugeProducer.mCurrentSlicedBucket->begin()->second->begin()->second.value_int()); + gaugeProducer.mCurrentSlicedBucket->begin()-> + second.front().mFields->begin()->second.value_int()); EXPECT_EQ(0UL, gaugeProducer.mPastBuckets.size()); vector> allData; @@ -165,17 +166,18 @@ TEST(GaugeMetricProducerTest, TestWithCondition) { EXPECT_EQ(1UL, gaugeProducer.mCurrentSlicedBucket->size()); EXPECT_EQ(110, - gaugeProducer.mCurrentSlicedBucket->begin()->second->begin()->second.value_int()); + gaugeProducer.mCurrentSlicedBucket->begin()-> + second.front().mFields->begin()->second.value_int()); EXPECT_EQ(1UL, gaugeProducer.mPastBuckets.size()); EXPECT_EQ(100, gaugeProducer.mPastBuckets.begin()->second.back() - .mGaugeFields->begin()->second.value_int()); + .mGaugeAtoms.front().mFields->begin()->second.value_int()); gaugeProducer.onConditionChanged(false, bucket2StartTimeNs + 10); gaugeProducer.flushIfNeededLocked(bucket3StartTimeNs + 10); EXPECT_EQ(1UL, gaugeProducer.mPastBuckets.size()); EXPECT_EQ(2UL, gaugeProducer.mPastBuckets.begin()->second.size()); EXPECT_EQ(110L, gaugeProducer.mPastBuckets.begin()->second.back() - .mGaugeFields->begin()->second.value_int()); + .mGaugeAtoms.front().mFields->begin()->second.value_int()); EXPECT_EQ(1UL, gaugeProducer.mPastBuckets.begin()->second.back().mBucketNum); } @@ -214,7 +216,8 @@ TEST(GaugeMetricProducerTest, TestAnomalyDetection) { gaugeProducer.onDataPulled({event1}); EXPECT_EQ(1UL, gaugeProducer.mCurrentSlicedBucket->size()); EXPECT_EQ(13L, - gaugeProducer.mCurrentSlicedBucket->begin()->second->begin()->second.value_int()); + gaugeProducer.mCurrentSlicedBucket->begin()-> + second.front().mFields->begin()->second.value_int()); EXPECT_EQ(anomalyTracker->getRefractoryPeriodEndsSec(DEFAULT_DIMENSION_KEY), 0U); std::shared_ptr event2 = @@ -226,7 +229,8 @@ TEST(GaugeMetricProducerTest, TestAnomalyDetection) { gaugeProducer.onDataPulled({event2}); EXPECT_EQ(1UL, gaugeProducer.mCurrentSlicedBucket->size()); EXPECT_EQ(15L, - gaugeProducer.mCurrentSlicedBucket->begin()->second->begin()->second.value_int()); + gaugeProducer.mCurrentSlicedBucket->begin()-> + second.front().mFields->begin()->second.value_int()); EXPECT_EQ(anomalyTracker->getRefractoryPeriodEndsSec(DEFAULT_DIMENSION_KEY), event2->GetTimestampNs() / NS_PER_SEC + refPeriodSec); @@ -239,7 +243,8 @@ TEST(GaugeMetricProducerTest, TestAnomalyDetection) { gaugeProducer.onDataPulled({event3}); EXPECT_EQ(1UL, gaugeProducer.mCurrentSlicedBucket->size()); EXPECT_EQ(26L, - gaugeProducer.mCurrentSlicedBucket->begin()->second->begin()->second.value_int()); + gaugeProducer.mCurrentSlicedBucket->begin()-> + second.front().mFields->begin()->second.value_int()); EXPECT_EQ(anomalyTracker->getRefractoryPeriodEndsSec(DEFAULT_DIMENSION_KEY), event2->GetTimestampNs() / NS_PER_SEC + refPeriodSec); @@ -250,7 +255,7 @@ TEST(GaugeMetricProducerTest, TestAnomalyDetection) { event4->init(); gaugeProducer.onDataPulled({event4}); EXPECT_EQ(1UL, gaugeProducer.mCurrentSlicedBucket->size()); - EXPECT_TRUE(gaugeProducer.mCurrentSlicedBucket->begin()->second->empty()); + EXPECT_TRUE(gaugeProducer.mCurrentSlicedBucket->begin()->second.front().mFields->empty()); } } // namespace statsd