Skip to content

Commit 8bee203

Browse files
committed
fix(executor): reject zero thread count
1 parent 342d25a commit 8bee203

14 files changed

Lines changed: 48 additions & 33 deletions

include/paimon/executor.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
#include <functional>
2121
#include <memory>
2222

23+
#include "paimon/result.h"
2324
#include "paimon/visibility.h"
2425

2526
namespace paimon {
@@ -34,7 +35,7 @@ PAIMON_EXPORT std::shared_ptr<Executor> GetGlobalDefaultExecutor();
3435
PAIMON_EXPORT std::unique_ptr<Executor> CreateDefaultExecutor();
3536

3637
/// Create a default implementation of executor with specified thread_count.
37-
PAIMON_EXPORT std::unique_ptr<Executor> CreateDefaultExecutor(uint32_t thread_count);
38+
PAIMON_EXPORT Result<std::unique_ptr<Executor>> CreateDefaultExecutor(uint32_t thread_count);
3839

3940
/// Interface class for defining basic operations of a task executor.
4041
///

src/paimon/common/executor/default_executor_test.cpp

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
#include "paimon/executor.h"
2929
#include "paimon/result.h"
3030
#include "paimon/status.h"
31+
#include "paimon/testing/utils/testharness.h"
3132

3233
namespace paimon::test {
3334

@@ -83,7 +84,7 @@ TEST(DefaultExecutorTest, TestViaWithException) {
8384
}
8485

8586
TEST(DefaultExecutorTest, TestShutdownNowDropsPendingTasks) {
86-
auto executor = CreateDefaultExecutor(/*thread_count=*/1);
87+
ASSERT_OK_AND_ASSIGN(auto executor, CreateDefaultExecutor(/*thread_count=*/1));
8788
std::atomic<bool> first_started = false;
8889
std::atomic<int32_t> executed_count = 0;
8990
std::promise<void> release_first_task;
@@ -112,7 +113,7 @@ TEST(DefaultExecutorTest, TestShutdownNowDropsPendingTasks) {
112113
}
113114

114115
TEST(DefaultExecutorTest, TestAddTaskAfterShutdownNowIgnored) {
115-
auto executor = CreateDefaultExecutor(/*thread_count=*/1);
116+
ASSERT_OK_AND_ASSIGN(auto executor, CreateDefaultExecutor(/*thread_count=*/1));
116117
std::atomic<int32_t> executed_count = 0;
117118

118119
executor->ShutdownNow();
@@ -122,4 +123,9 @@ TEST(DefaultExecutorTest, TestAddTaskAfterShutdownNowIgnored) {
122123
ASSERT_EQ(executed_count.load(), 0);
123124
}
124125

126+
TEST(DefaultExecutorTest, TestCreateWithZeroThreadCount) {
127+
ASSERT_NOK_WITH_MSG(CreateDefaultExecutor(/*thread_count=*/0),
128+
"default executor thread count should be greater than 0");
129+
}
130+
125131
} // namespace paimon::test

src/paimon/common/executor/executor.cpp

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

1717
#include "paimon/executor.h"
1818

19+
#include <cassert>
1920
#include <condition_variable>
2021
#include <functional>
2122
#include <mutex>
@@ -49,8 +50,8 @@ class DefaultExecutor : public Executor {
4950
int32_t active_tasks_ = 0;
5051
};
5152

52-
DefaultExecutor::DefaultExecutor(uint32_t thread_count)
53-
: thread_count_(thread_count == 0 ? 1 : thread_count) {
53+
DefaultExecutor::DefaultExecutor(uint32_t thread_count) : thread_count_(thread_count) {
54+
assert(thread_count > 0);
5455
for (uint32_t i = 0; i < thread_count_; ++i) {
5556
workers_.emplace_back(&DefaultExecutor::WorkerThread, this);
5657
}
@@ -136,10 +137,13 @@ PAIMON_EXPORT std::shared_ptr<Executor> GetGlobalDefaultExecutor() {
136137
}
137138

138139
PAIMON_EXPORT std::unique_ptr<Executor> CreateDefaultExecutor() {
139-
return CreateDefaultExecutor(DEFAULT_EXECUTOR_THREAD_COUNT);
140+
return std::make_unique<DefaultExecutor>(DEFAULT_EXECUTOR_THREAD_COUNT);
140141
}
141142

142-
PAIMON_EXPORT std::unique_ptr<Executor> CreateDefaultExecutor(uint32_t thread_count) {
143+
PAIMON_EXPORT Result<std::unique_ptr<Executor>> CreateDefaultExecutor(uint32_t thread_count) {
144+
if (thread_count == 0) {
145+
return Status::Invalid("default executor thread count should be greater than 0");
146+
}
143147
return std::make_unique<DefaultExecutor>(thread_count);
144148
}
145149

src/paimon/common/file_index/bitmap/apply_bitmap_index_batch_reader_test.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,7 +54,7 @@ class ApplyBitmapIndexBatchReaderTest : public ::testing::Test,
5454

5555
pool_ = GetDefaultPool();
5656
fs_ = std::make_shared<MockFileSystem>();
57-
executor_ = CreateDefaultExecutor(/*thread_count=*/2);
57+
ASSERT_OK_AND_ASSIGN(executor_, CreateDefaultExecutor(/*thread_count=*/2));
5858
}
5959
void TearDown() override {}
6060

src/paimon/common/fs/file_system_test.cpp

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1086,7 +1086,7 @@ TEST_P(FileSystemTest, TestMkdir2) {
10861086
TEST_P(FileSystemTest, TestMkdirMultiThreadWithSameNonExistParentDir) {
10871087
uint32_t runs_count = 10;
10881088
uint32_t thread_count = 10;
1089-
auto executor = CreateDefaultExecutor(thread_count);
1089+
ASSERT_OK_AND_ASSIGN(auto executor, CreateDefaultExecutor(thread_count));
10901090

10911091
for (uint32_t i = 0; i < runs_count; i++) {
10921092
std::string uuid;
@@ -1111,7 +1111,7 @@ TEST_P(FileSystemTest, TestMkdirMultiThreadWithSameNonExistParentDir) {
11111111
TEST_P(FileSystemTest, TestMkdirMultiThreadWithSameName) {
11121112
uint32_t runs_count = 10;
11131113
uint32_t thread_count = 10;
1114-
auto executor = CreateDefaultExecutor(thread_count);
1114+
ASSERT_OK_AND_ASSIGN(auto executor, CreateDefaultExecutor(thread_count));
11151115

11161116
for (uint32_t i = 0; i < runs_count; i++) {
11171117
std::string uuid;
@@ -1133,7 +1133,7 @@ TEST_P(FileSystemTest, TestMkdirMultiThreadWithSameName) {
11331133
TEST_P(FileSystemTest, TestMkdirMultiThreadWithSameNameWithRelativePath) {
11341134
uint32_t runs_count = 10;
11351135
uint32_t thread_count = 10;
1136-
auto executor = CreateDefaultExecutor(thread_count);
1136+
ASSERT_OK_AND_ASSIGN(auto executor, CreateDefaultExecutor(thread_count));
11371137

11381138
for (uint32_t i = 0; i < runs_count; i++) {
11391139
std::string uuid;

src/paimon/common/logging/logging_test.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@
2424
#include "paimon/testing/utils/testharness.h"
2525
namespace paimon::test {
2626
TEST(LoggerTest, TestMultiThreadGetLogger) {
27-
auto executor = CreateDefaultExecutor(/*thread_count=*/4);
27+
ASSERT_OK_AND_ASSIGN(auto executor, CreateDefaultExecutor(/*thread_count=*/4));
2828
auto get_logger = []() {
2929
auto logger = Logger::GetLogger("my_log");
3030
ASSERT_TRUE(logger);

src/paimon/common/reader/prefetch_file_batch_reader_impl_test.cpp

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -127,7 +127,7 @@ class PrefetchFileBatchReaderImplTest : public ::testing::Test,
127127
data_type_ = arrow::struct_(fields_);
128128
mock_fs_ = std::make_shared<MockFileSystem>();
129129
local_fs_ = std::make_shared<LocalFileSystem>();
130-
executor_ = CreateDefaultExecutor(/*thread_count=*/2);
130+
ASSERT_OK_AND_ASSIGN(executor_, CreateDefaultExecutor(/*thread_count=*/2));
131131
dir_ = ::paimon::test::UniqueTestDirectory::Create();
132132
ASSERT_TRUE(dir_);
133133
}
@@ -198,14 +198,16 @@ class PrefetchFileBatchReaderImplTest : public ::testing::Test,
198198
EXPECT_OK_AND_ASSIGN(std::unique_ptr<FileFormat> file_format,
199199
FileFormatFactory::Get(file_format_str, {}));
200200
EXPECT_OK_AND_ASSIGN(auto reader_builder, file_format->CreateReaderBuilder(batch_size));
201+
EXPECT_OK_AND_ASSIGN(std::shared_ptr<Executor> executor,
202+
CreateDefaultExecutor(prefetch_max_parallel_num - 1));
201203
EXPECT_OK_AND_ASSIGN(
202204
std::unique_ptr<PrefetchFileBatchReaderImpl> reader,
203205
PrefetchFileBatchReaderImpl::Create(
204206
PathUtil::JoinPath(dir_->Str(), "file." + file_format->Identifier()),
205207
reader_builder.get(), local_fs_, prefetch_max_parallel_num, batch_size,
206208
prefetch_max_parallel_num * 2, /*enable_adaptive_prefetch_strategy=*/false,
207-
CreateDefaultExecutor(prefetch_max_parallel_num - 1),
208-
/*initialize_read_ranges=*/false, cache_mode, CacheConfig(), GetDefaultPool()));
209+
executor, /*initialize_read_ranges=*/false, cache_mode, CacheConfig(),
210+
GetDefaultPool()));
209211
std::unique_ptr<ArrowSchema> c_schema = std::make_unique<ArrowSchema>();
210212
auto arrow_status = arrow::ExportSchema(*read_schema, c_schema.get());
211213
EXPECT_TRUE(arrow_status.ok());

src/paimon/core/deletionvectors/apply_deletion_vector_batch_reader_test.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,7 @@ class ApplyDeletionVectorBatchReaderTest : public ::testing::Test,
4949

5050
pool_ = GetDefaultPool();
5151
fs_ = std::make_shared<MockFileSystem>();
52-
executor_ = CreateDefaultExecutor(/*thread_count=*/2);
52+
ASSERT_OK_AND_ASSIGN(executor_, CreateDefaultExecutor(/*thread_count=*/2));
5353
}
5454
void TearDown() override {}
5555

src/paimon/core/global_index/global_index_scan_impl.cpp

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -107,10 +107,8 @@ Result<std::unique_ptr<GlobalIndexScanImpl>> GlobalIndexScanImpl::Create(
107107
uint32_t cpu_count = std::thread::hardware_concurrency();
108108
thread_num = cpu_count > 0 ? static_cast<int32_t>(cpu_count) : 1;
109109
}
110-
if (thread_num.value() <= 0) {
111-
thread_num = 1;
112-
}
113-
final_executor = CreateDefaultExecutor(static_cast<uint32_t>(thread_num.value()));
110+
PAIMON_ASSIGN_OR_RAISE(final_executor,
111+
CreateDefaultExecutor(static_cast<uint32_t>(thread_num.value())));
114112
}
115113
return std::unique_ptr<GlobalIndexScanImpl>(new GlobalIndexScanImpl(
116114
table_schema, options, path_factory, std::move(index_metas), final_executor, pool));

src/paimon/core/operation/abstract_file_store_write.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -74,7 +74,7 @@ AbstractFileStoreWrite::AbstractFileStoreWrite(
7474
dv_maintainer_factory_(dv_maintainer_factory),
7575
io_manager_(io_manager),
7676
options_(options),
77-
compact_executor_(CreateDefaultExecutor(4)),
77+
compact_executor_(CreateDefaultExecutor()),
7878
compaction_metrics_(std::make_shared<CompactionMetrics>()),
7979
ignore_previous_files_(ignore_previous_files),
8080
is_streaming_mode_(is_streaming_mode),

0 commit comments

Comments
 (0)