Thread pool design point:
All tasks submitted directly to the thread pool enter a FIFO queue and are
dispatched to a worker thread when one becomes free. Tasks may also be
submitted via ThreadPoolTokens. The token wait() and shutdown() functions
can then be used to block on logical groups of tasks.
A token operates in one of two ExecutionModes, determined at token
construction time:
1. SERIAL: submitted tasks are run one at a time.
2. CONCURRENT: submitted tasks may be run in parallel.
This isn't unlike submitted without a token, but the logical grouping that tokens
impart can be useful when a pool is shared by many contexts (e.g. to
safely shut down one context, to derive context-specific metrics, etc.).
Tasks submitted without a token or via ExecutionMode::CONCURRENT tokens are
processed in FIFO order. On the other hand, ExecutionMode::SERIAL tokens are
processed in a round-robin fashion, one task at a time. This prevents them
from starving one another. However, tokenless (and CONCURRENT token-based)
tasks can starve SERIAL token-based tasks.
Thread design point:
1. It is a thin wrapper around pthread that can register itself with the singleton ThreadMgr
(a private class implemented in thread.cpp entirely, which tracks all live threads so
that they may be monitored via the debug webpages). This class has a limited subset of
boost::thread's API. Construction is almost the same, but clients must supply a
category and a name for each thread so that they can be identified in the debug web
UI. Otherwise, join() is the only supported method from boost::thread.
2. Each Thread object knows its operating system thread ID (TID), which can be used to
attach debuggers to specific threads, to retrieve resource-usage statistics from the
operating system, and to assign threads to resource control groups.
3. Threads are shared objects, but in a degenerate way. They may only have
up to two referents: the caller that created the thread (parent), and
the thread itself (child). Moreover, the only two methods to mutate state
(join() and the destructor) are constrained: the child may not join() on
itself, and the destructor is only run when there's one referent left.
These constraints allow us to access thread internals without any locks.
78 lines
2.5 KiB
C++
78 lines
2.5 KiB
C++
// Licensed to the Apache Software Foundation (ASF) under one
|
|
// or more contributor license agreements. See the NOTICE file
|
|
// distributed with this work for additional information
|
|
// regarding copyright ownership. The ASF licenses this file
|
|
// to you under the Apache License, Version 2.0 (the
|
|
// "License"); you may not use this file except in compliance
|
|
// with the License. You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing,
|
|
// software distributed under the License is distributed on an
|
|
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
|
// KIND, either express or implied. See the License for the
|
|
// specific language governing permissions and limitations
|
|
// under the License.
|
|
|
|
#include <functional>
|
|
#include <gtest/gtest.h>
|
|
|
|
#include "gutil/ref_counted.h"
|
|
#include "util/countdown_latch.h"
|
|
#include "util/monotime.h"
|
|
#include "util/thread.h"
|
|
#include "util/threadpool.h"
|
|
|
|
namespace doris {
|
|
|
|
static void decrement_latch(CountDownLatch* latch, int amount) {
|
|
if (amount == 1) {
|
|
latch->count_down();
|
|
return;
|
|
}
|
|
latch->count_down(amount);
|
|
}
|
|
|
|
// Tests that we can decrement the latch by arbitrary amounts, as well
|
|
// as 1 by one.
|
|
TEST(TestCountDownLatch, TestLatch) {
|
|
|
|
std::unique_ptr<ThreadPool> pool;
|
|
ASSERT_TRUE(ThreadPoolBuilder("cdl-test").set_max_threads(1).build(&pool).ok());
|
|
|
|
CountDownLatch latch(1000);
|
|
|
|
// Decrement the count by 1 in another thread, this should not fire the
|
|
// latch.
|
|
ASSERT_TRUE(pool->submit_func(std::bind(decrement_latch, &latch, 1)).ok());
|
|
ASSERT_FALSE(latch.wait_for(MonoDelta::FromMilliseconds(200)));
|
|
ASSERT_EQ(999, latch.count());
|
|
|
|
// Now decrement by 1000 this should decrement to 0 and fire the latch
|
|
// (even though 1000 is one more than the current count).
|
|
ASSERT_TRUE(pool->submit_func(std::bind(decrement_latch, &latch, 1000)).ok());
|
|
latch.wait();
|
|
ASSERT_EQ(0, latch.count());
|
|
}
|
|
|
|
// Test that resetting to zero while there are waiters lets the waiters
|
|
// continue.
|
|
TEST(TestCountDownLatch, TestResetToZero) {
|
|
CountDownLatch cdl(100);
|
|
scoped_refptr<Thread> t;
|
|
ASSERT_TRUE(Thread::create("test", "cdl-test", &CountDownLatch::wait, &cdl, &t).ok());
|
|
|
|
// Sleep for a bit until it's likely the other thread is waiting on the latch.
|
|
SleepFor(MonoDelta::FromMilliseconds(10));
|
|
cdl.reset(0);
|
|
t->join();
|
|
}
|
|
|
|
} // namespace doris
|
|
|
|
int main(int argc, char* argv[]) {
|
|
::testing::InitGoogleTest(&argc, argv);
|
|
return RUN_ALL_TESTS();
|
|
}
|