Skip to content

Commit 0c0610e

Browse files
committed
add tests for executor
1 parent ee6159e commit 0c0610e

1 file changed

Lines changed: 54 additions & 0 deletions

File tree

src/paimon/common/executor/default_executor_test.cpp

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -123,6 +123,60 @@ TEST(DefaultExecutorTest, TestAddTaskAfterShutdownNowIgnored) {
123123
ASSERT_EQ(executed_count.load(), 0);
124124
}
125125

126+
TEST(DefaultExecutorTest, TestAddTaskFromMultipleThreads) {
127+
ASSERT_OK_AND_ASSIGN(auto executor, CreateDefaultExecutor(/*thread_count=*/4));
128+
129+
constexpr int32_t kSubmitterCount = 8;
130+
constexpr int32_t kTaskCountPerSubmitter = 64;
131+
constexpr int32_t kTotalTaskCount = kSubmitterCount * kTaskCountPerSubmitter;
132+
133+
std::vector<std::atomic<int32_t>> executed_slots(kTotalTaskCount);
134+
for (auto& executed_slot : executed_slots) {
135+
executed_slot.store(0);
136+
}
137+
std::vector<std::promise<void>> task_promises(kTotalTaskCount);
138+
std::vector<std::future<void>> task_futures;
139+
task_futures.reserve(kTotalTaskCount);
140+
for (auto& task_promise : task_promises) {
141+
task_futures.push_back(task_promise.get_future());
142+
}
143+
std::atomic<int32_t> ready_submitter_count = 0;
144+
std::atomic<int32_t> executed_count = 0;
145+
std::promise<void> start_signal;
146+
std::shared_future<void> start_future = start_signal.get_future().share();
147+
std::vector<std::thread> submitters;
148+
submitters.reserve(kSubmitterCount);
149+
150+
for (int32_t submitter_index = 0; submitter_index < kSubmitterCount; ++submitter_index) {
151+
submitters.emplace_back([&, submitter_index]() {
152+
++ready_submitter_count;
153+
start_future.wait();
154+
for (int32_t task_index = 0; task_index < kTaskCountPerSubmitter; ++task_index) {
155+
const int32_t slot_index = submitter_index * kTaskCountPerSubmitter + task_index;
156+
executor->Add([&, slot_index]() {
157+
++executed_slots[slot_index];
158+
++executed_count;
159+
task_promises[slot_index].set_value();
160+
});
161+
}
162+
});
163+
}
164+
165+
while (ready_submitter_count.load() < kSubmitterCount) {
166+
std::this_thread::yield();
167+
}
168+
start_signal.set_value();
169+
for (auto& submitter : submitters) {
170+
submitter.join();
171+
}
172+
Wait(task_futures);
173+
174+
ASSERT_EQ(kTotalTaskCount, executed_count.load());
175+
for (const auto& executed_slot : executed_slots) {
176+
ASSERT_EQ(1, executed_slot.load());
177+
}
178+
}
179+
126180
TEST(DefaultExecutorTest, TestCreateWithZeroThreadCount) {
127181
ASSERT_NOK_WITH_MSG(CreateDefaultExecutor(/*thread_count=*/0),
128182
"default executor thread count should be greater than 0");

0 commit comments

Comments
 (0)