mirror of
https://github.com/wassname/ray.git
synced 2026-07-25 13:30:52 +08:00
[Stats] Fix harvestor threads + Fix flaky stats shutdown. (#9745)
This commit is contained in:
@@ -109,6 +109,7 @@ static inline void Shutdown() {
|
||||
return;
|
||||
}
|
||||
metrics_io_service_pool->Stop();
|
||||
opencensus::stats::DeltaProducer::Get()->Shutdown();
|
||||
opencensus::stats::StatsExporter::Shutdown();
|
||||
metrics_io_service_pool = nullptr;
|
||||
exporter = nullptr;
|
||||
|
||||
@@ -119,6 +119,62 @@ TEST_F(StatsTest, InitializationTest) {
|
||||
ASSERT_TRUE(new_first_tag.second == test_tag_value_that_shouldnt_be_applied);
|
||||
}
|
||||
|
||||
TEST_F(StatsTest, MultiThreadedInitializationTest) {
|
||||
// Make sure stats module is thread-safe.
|
||||
// Shutdown the stats module first.
|
||||
ray::stats::Shutdown();
|
||||
// Spawn 10 threads that init and shutdown again and again.
|
||||
// The test will have memory corruption if it doesn't work as expected.
|
||||
const stats::TagsType global_tags = {{stats::LanguageKey, "CPP"},
|
||||
{stats::WorkerPidKey, "1000"}};
|
||||
std::vector<std::thread> threads;
|
||||
for (int i = 0; i < 5; i++) {
|
||||
threads.emplace_back([global_tags]() {
|
||||
for (int i = 0; i < 5; i++) {
|
||||
std::shared_ptr<stats::MetricExporterClient> exporter(
|
||||
new stats::StdoutExporterClient());
|
||||
unsigned int upper_bound = 100;
|
||||
unsigned int init_or_shutdown = (rand() % upper_bound);
|
||||
if (init_or_shutdown >= (upper_bound / 2)) {
|
||||
ray::stats::Init(global_tags, MetricsAgentPort, exporter);
|
||||
} else {
|
||||
ray::stats::Shutdown();
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
for (auto &thread : threads) {
|
||||
thread.join();
|
||||
}
|
||||
ray::stats::Shutdown();
|
||||
ASSERT_FALSE(ray::stats::StatsConfig::instance().IsInitialized());
|
||||
std::shared_ptr<stats::MetricExporterClient> exporter(
|
||||
new stats::StdoutExporterClient());
|
||||
ray::stats::Init(global_tags, MetricsAgentPort, exporter);
|
||||
ASSERT_TRUE(ray::stats::StatsConfig::instance().IsInitialized());
|
||||
}
|
||||
|
||||
TEST_F(StatsTest, TestShutdownTakesLongTime) {
|
||||
// Make sure it doesn't take long time to shutdown when harvestor / export interval is
|
||||
// large.
|
||||
ray::stats::Shutdown();
|
||||
// Spawn 10 threads that init and shutdown again and again.
|
||||
// The test will have memory corruption if it doesn't work as expected.
|
||||
const stats::TagsType global_tags = {{stats::LanguageKey, "CPP"},
|
||||
{stats::WorkerPidKey, "1000"}};
|
||||
std::shared_ptr<stats::MetricExporterClient> exporter(
|
||||
new stats::StdoutExporterClient());
|
||||
|
||||
// Flush interval is 30 seconds. Shutdown should not take 30 seconds in this case.
|
||||
uint32_t kReportFlushInterval = 30000;
|
||||
absl::Duration report_interval = absl::Milliseconds(kReportFlushInterval);
|
||||
absl::Duration harvest_interval = absl::Milliseconds(kReportFlushInterval);
|
||||
ray::stats::StatsConfig::instance().SetReportInterval(report_interval);
|
||||
ray::stats::StatsConfig::instance().SetHarvestInterval(harvest_interval);
|
||||
ray::stats::Init(global_tags, MetricsAgentPort, exporter);
|
||||
ray::stats::Shutdown();
|
||||
}
|
||||
|
||||
} // namespace ray
|
||||
|
||||
int main(int argc, char **argv) {
|
||||
|
||||
Reference in New Issue
Block a user