diff --git a/docs/config/client.en.md b/docs/config/client.en.md index a72a3e41..3673bd86 100644 --- a/docs/config/client.en.md +++ b/docs/config/client.en.md @@ -9,6 +9,7 @@ These parameters apply only to Vitastor clients (QEMU, fio, NBD and so on) and affect their interaction with the cluster. +- [client_iothread_count](#client_iothread_count) - [client_retry_interval](#client_retry_interval) - [client_eio_retry_interval](#client_eio_retry_interval) - [client_retry_enospc](#client_retry_enospc) @@ -23,6 +24,23 @@ affect their interaction with the cluster. - [nbd_max_part](#nbd_max_part) - [osd_nearfull_ratio](#osd_nearfull_ratio) +## client_iothread_count + +- Type: integer +- Default: 0 + +Number of separate threads for handling TCP network I/O at client library +side. Enabling 4 threads usually allows to increase peak performance of each +client from approx. 2-3 to 7-8 GByte/s linear read/write and from approx. +100-150 to 400 thousand iops, but at the same time it increases latency. +Latency increase depends on CPU: with CPU power saving disabled latency +only increases by ~10 us (equivalent to Q=1 iops decrease from 10500 to 9500), +with CPU power saving enabled it may be as high as 500 us (equivalent to Q=1 +iops decrease from 2000 to 1000). RDMA isn't affected by this option. + +It's recommended to enable client I/O threads if you don't use RDMA and want +to increase peak client performance. + ## client_retry_interval - Type: milliseconds diff --git a/docs/config/client.ru.md b/docs/config/client.ru.md index 2f75eacc..e941e642 100644 --- a/docs/config/client.ru.md +++ b/docs/config/client.ru.md @@ -9,6 +9,7 @@ Данные параметры применяются только к клиентам Vitastor (QEMU, fio, NBD и т.п.) и затрагивают логику их работы с кластером. +- [client_iothread_count](#client_iothread_count) - [client_retry_interval](#client_retry_interval) - [client_eio_retry_interval](#client_eio_retry_interval) - [client_retry_enospc](#client_retry_enospc) @@ -23,6 +24,24 @@ - [nbd_max_part](#nbd_max_part) - [osd_nearfull_ratio](#osd_nearfull_ratio) +## client_iothread_count + +- Тип: целое число +- Значение по умолчанию: 0 + +Число отдельных потоков для обработки ввода-вывода через TCP сеть на стороне +клиентской библиотеки. Включение 4 потоков обычно позволяет поднять пиковую +производительность каждого клиента примерно с 2-3 до 7-8 Гбайт/с линейного +чтения/записи и примерно с 100-150 до 400 тысяч операций ввода-вывода в +секунду, но ухудшает задержку. Увеличение задержки зависит от процессора: +при отключённом энергосбережении CPU это всего ~10 микросекунд (равносильно +падению iops с Q=1 с 10500 до 9500), а при включённом это может быть +и 500 микросекунд (равносильно падению iops с Q=1 с 2000 до 1000). На работу +RDMA данная опция не влияет. + +Рекомендуется включать клиентские потоки ввода-вывода, если вы не используете +RDMA и хотите повысить пиковую производительность клиентов. + ## client_retry_interval - Тип: миллисекунды diff --git a/docs/config/osd.en.md b/docs/config/osd.en.md index 165a79df..9fda3ff9 100644 --- a/docs/config/osd.en.md +++ b/docs/config/osd.en.md @@ -10,6 +10,7 @@ These parameters only apply to OSDs, are not fixed at the moment of OSD drive initialization and can be changed - either with an OSD restart or, for some of them, even without restarting by updating configuration in etcd. +- [osd_iothread_count](#osd_iothread_count) - [etcd_report_interval](#etcd_report_interval) - [etcd_stats_interval](#etcd_stats_interval) - [run_primary](#run_primary) @@ -61,6 +62,18 @@ them, even without restarting by updating configuration in etcd. - [recovery_tune_sleep_min_us](#recovery_tune_sleep_min_us) - [recovery_tune_sleep_cutoff_us](#recovery_tune_sleep_cutoff_us) +## osd_iothread_count + +- Type: integer +- Default: 0 + +TCP network I/O thread count for OSD. When non-zero, a single OSD process +may handle more TCP I/O, but at a cost of increased latency because thread +switching overhead occurs. RDMA isn't affected by this option. + +Because of latency, instead of enabling OSD I/O threads it's recommended to +just create multiple OSDs per disk, or use RDMA. + ## etcd_report_interval - Type: seconds diff --git a/docs/config/osd.ru.md b/docs/config/osd.ru.md index 66456088..a70a743b 100644 --- a/docs/config/osd.ru.md +++ b/docs/config/osd.ru.md @@ -11,6 +11,7 @@ момент с помощью перезапуска OSD, а некоторые и без перезапуска, с помощью изменения конфигурации в etcd. +- [osd_iothread_count](#osd_iothread_count) - [etcd_report_interval](#etcd_report_interval) - [etcd_stats_interval](#etcd_stats_interval) - [run_primary](#run_primary) @@ -62,6 +63,19 @@ - [recovery_tune_sleep_min_us](#recovery_tune_sleep_min_us) - [recovery_tune_sleep_cutoff_us](#recovery_tune_sleep_cutoff_us) +## osd_iothread_count + +- Тип: целое число +- Значение по умолчанию: 0 + +Число отдельных потоков для обработки ввода-вывода через TCP-сеть на +стороне OSD. Включение опции позволяет каждому отдельному OSD передавать +по сети больше данных, но ухудшает задержку из-за накладных расходов +переключения потоков. На работу RDMA опция не влияет. + +Из-за задержек вместо включения потоков ввода-вывода OSD рекомендуется +просто создавать по несколько OSD на каждом диске, или использовать RDMA. + ## etcd_report_interval - Тип: секунды diff --git a/docs/config/src/client.yml b/docs/config/src/client.yml index 842fa1bf..4d46a7cf 100644 --- a/docs/config/src/client.yml +++ b/docs/config/src/client.yml @@ -1,3 +1,32 @@ +- name: client_iothread_count + type: int + default: 0 + online: false + info: | + Number of separate threads for handling TCP network I/O at client library + side. Enabling 4 threads usually allows to increase peak performance of each + client from approx. 2-3 to 7-8 GByte/s linear read/write and from approx. + 100-150 to 400 thousand iops, but at the same time it increases latency. + Latency increase depends on CPU: with CPU power saving disabled latency + only increases by ~10 us (equivalent to Q=1 iops decrease from 10500 to 9500), + with CPU power saving enabled it may be as high as 500 us (equivalent to Q=1 + iops decrease from 2000 to 1000). RDMA isn't affected by this option. + + It's recommended to enable client I/O threads if you don't use RDMA and want + to increase peak client performance. + info_ru: | + Число отдельных потоков для обработки ввода-вывода через TCP сеть на стороне + клиентской библиотеки. Включение 4 потоков обычно позволяет поднять пиковую + производительность каждого клиента примерно с 2-3 до 7-8 Гбайт/с линейного + чтения/записи и примерно с 100-150 до 400 тысяч операций ввода-вывода в + секунду, но ухудшает задержку. Увеличение задержки зависит от процессора: + при отключённом энергосбережении CPU это всего ~10 микросекунд (равносильно + падению iops с Q=1 с 10500 до 9500), а при включённом это может быть + и 500 микросекунд (равносильно падению iops с Q=1 с 2000 до 1000). На работу + RDMA данная опция не влияет. + + Рекомендуется включать клиентские потоки ввода-вывода, если вы не используете + RDMA и хотите повысить пиковую производительность клиентов. - name: client_retry_interval type: ms min: 10 diff --git a/docs/config/src/osd.yml b/docs/config/src/osd.yml index 474ed8bf..93d0e266 100644 --- a/docs/config/src/osd.yml +++ b/docs/config/src/osd.yml @@ -1,3 +1,21 @@ +- name: osd_iothread_count + type: int + default: 0 + info: | + TCP network I/O thread count for OSD. When non-zero, a single OSD process + may handle more TCP I/O, but at a cost of increased latency because thread + switching overhead occurs. RDMA isn't affected by this option. + + Because of latency, instead of enabling OSD I/O threads it's recommended to + just create multiple OSDs per disk, or use RDMA. + info_ru: | + Число отдельных потоков для обработки ввода-вывода через TCP-сеть на + стороне OSD. Включение опции позволяет каждому отдельному OSD передавать + по сети больше данных, но ухудшает задержку из-за накладных расходов + переключения потоков. На работу RDMA опция не влияет. + + Из-за задержек вместо включения потоков ввода-вывода OSD рекомендуется + просто создавать по несколько OSD на каждом диске, или использовать RDMA. - name: etcd_report_interval type: sec default: 5 diff --git a/src/client/CMakeLists.txt b/src/client/CMakeLists.txt index 701f48bc..68512b03 100644 --- a/src/client/CMakeLists.txt +++ b/src/client/CMakeLists.txt @@ -12,6 +12,7 @@ add_library(vitastor_common STATIC msgr_stop.cpp msgr_op.cpp msgr_send.cpp msgr_receive.cpp ../util/ringloop.cpp ../../json11/json11.cpp http_client.cpp osd_ops.cpp pg_states.cpp ../util/timerfd_manager.cpp ../util/str_util.cpp ${MSGR_RDMA} ) +target_link_libraries(vitastor_common pthread) target_compile_options(vitastor_common PUBLIC -fPIC) # libvitastor_client.so diff --git a/src/client/messenger.cpp b/src/client/messenger.cpp index 74d0d2d2..075a8f7a 100644 --- a/src/client/messenger.cpp +++ b/src/client/messenger.cpp @@ -15,6 +15,106 @@ #include "msgr_rdma.h" #endif +#include + +msgr_iothread_t::msgr_iothread_t(): + ring(RINGLOOP_DEFAULT_SIZE, true), + thread(&msgr_iothread_t::run, this) +{ + eventfd = ring.register_eventfd(); + if (eventfd < 0) + { + throw std::runtime_error(std::string("failed to register eventfd: ") + strerror(-eventfd)); + } +} + +msgr_iothread_t::~msgr_iothread_t() +{ + stop(); +} + +void msgr_iothread_t::add_sqe(io_uring_sqe & sqe) +{ + mu.lock(); + queue.push_back((iothread_sqe_t){ .sqe = sqe, .data = std::move(*(ring_data_t*)sqe.user_data) }); + if (queue.size() == 1) + { + cond.notify_all(); + } + mu.unlock(); +} + +void msgr_iothread_t::stop() +{ + mu.lock(); + if (stopped) + { + mu.unlock(); + return; + } + stopped = true; + if (outer_loop_data) + { + outer_loop_data->callback = [](ring_data_t*){}; + } + cond.notify_all(); + close(eventfd); + mu.unlock(); + thread.join(); +} + +void msgr_iothread_t::add_to_ringloop(ring_loop_t *outer_loop) +{ + assert(!this->outer_loop || this->outer_loop == outer_loop); + io_uring_sqe *sqe = outer_loop->get_sqe(); + assert(sqe != NULL); + this->outer_loop = outer_loop; + this->outer_loop_data = ((ring_data_t*)sqe->user_data); + my_uring_prep_poll_add(sqe, eventfd, POLLIN); + outer_loop_data->callback = [this](ring_data_t *data) + { + if (data->res < 0) + { + throw std::runtime_error(std::string("eventfd poll failed: ") + strerror(-data->res)); + } + outer_loop_data = NULL; + if (stopped) + { + return; + } + add_to_ringloop(this->outer_loop); + ring.loop(); + }; +} + +void msgr_iothread_t::run() +{ + while (true) + { + { + std::unique_lock lk(mu); + while (!stopped && !queue.size()) + cond.wait(lk); + if (stopped) + return; + int i = 0; + for (; i < queue.size(); i++) + { + io_uring_sqe *sqe = ring.get_sqe(); + if (!sqe) + break; + ring_data_t *data = ((ring_data_t*)sqe->user_data); + *data = std::move(queue[i].data); + *sqe = queue[i].sqe; + sqe->user_data = (uint64_t)data; + } + queue.erase(queue.begin(), queue.begin()+i); + } + // We only want to offload sendmsg/recvmsg. Callbacks will be called in main thread + ring.submit(); + } +} + void osd_messenger_t::init() { #ifdef WITH_RDMA @@ -43,6 +143,15 @@ void osd_messenger_t::init() } } #endif + if (ringloop && iothread_count > 0) + { + for (int i = 0; i < iothread_count; i++) + { + auto iot = new msgr_iothread_t(); + iothreads.push_back(iot); + iot->add_to_ringloop(ringloop); + } + } keepalive_timer_id = tfd->set_timer(1000, true, [this](int) { auto cl_it = clients.begin(); @@ -129,6 +238,14 @@ osd_messenger_t::~osd_messenger_t() { stop_client(clients.begin()->first, true, true); } + if (iothreads.size()) + { + for (auto iot: iothreads) + { + delete iot; + } + iothreads.clear(); + } #ifdef WITH_RDMA if (rdma_context) { @@ -165,6 +282,10 @@ void osd_messenger_t::parse_config(const json11::Json & config) this->rdma_max_msg = 129*1024; this->rdma_odp = config["rdma_odp"].bool_value(); #endif + if (!osd_num) + this->iothread_count = (uint32_t)config["client_iothread_count"].uint64_value(); + else + this->iothread_count = (uint32_t)config["osd_iothread_count"].uint64_value(); this->receive_buffer_size = (uint32_t)config["tcp_header_buffer_size"].uint64_value(); if (!this->receive_buffer_size || this->receive_buffer_size > 1024*1024*1024) this->receive_buffer_size = 65536; diff --git a/src/client/messenger.h b/src/client/messenger.h index 6c44fbfe..3d447cdb 100644 --- a/src/client/messenger.h +++ b/src/client/messenger.h @@ -111,6 +111,44 @@ struct osd_op_stats_t uint64_t subop_stat_count[OSD_OP_MAX+1] = { 0 }; }; +#include +#include +#include + +#ifdef __MOCK__ +class msgr_iothread_t; +#else +struct iothread_sqe_t +{ + io_uring_sqe sqe; + ring_data_t data; +}; + +class msgr_iothread_t +{ +protected: + ring_loop_t ring; + ring_loop_t *outer_loop = NULL; + ring_data_t *outer_loop_data = NULL; + int eventfd = -1; + bool stopped = false; + std::mutex mu; + std::condition_variable cond; + std::vector queue; + std::thread thread; + + void run(); +public: + + msgr_iothread_t(); + ~msgr_iothread_t(); + + void add_sqe(io_uring_sqe & sqe); + void stop(); + void add_to_ringloop(ring_loop_t *outer_loop); +}; +#endif + struct osd_messenger_t { protected: @@ -123,6 +161,7 @@ protected: int osd_ping_timeout = 0; int log_level = 0; bool use_sync_send_recv = false; + int iothread_count = 0; #ifdef WITH_RDMA bool use_rdma = true; @@ -134,6 +173,7 @@ protected: bool rdma_odp = false; #endif + std::vector iothreads; std::vector read_ready_clients; std::vector write_ready_clients; // We don't use ringloop->set_immediate here because we may have no ringloop in client :) diff --git a/src/client/msgr_receive.cpp b/src/client/msgr_receive.cpp index 8534e091..3cc613b6 100644 --- a/src/client/msgr_receive.cpp +++ b/src/client/msgr_receive.cpp @@ -30,7 +30,11 @@ void osd_messenger_t::read_requests() cl->refs++; if (ringloop && !use_sync_send_recv) { - io_uring_sqe* sqe = ringloop->get_sqe(); + auto iothread = iothreads.size() ? iothreads[peer_fd % iothreads.size()] : NULL; + io_uring_sqe sqe_local; + ring_data_t data_local; + sqe_local.user_data = (uint64_t)&data_local; + io_uring_sqe* sqe = (iothread ? &sqe_local : ringloop->get_sqe()); if (!sqe) { cl->read_msg.msg_iovlen = 0; @@ -40,6 +44,10 @@ void osd_messenger_t::read_requests() ring_data_t* data = ((ring_data_t*)sqe->user_data); data->callback = [this, cl](ring_data_t *data) { handle_read(data->res, cl); }; my_uring_prep_recvmsg(sqe, peer_fd, &cl->read_msg, 0); + if (iothread) + { + iothread->add_sqe(sqe_local); + } } else { diff --git a/src/client/msgr_send.cpp b/src/client/msgr_send.cpp index 7374a1e7..adf70540 100644 --- a/src/client/msgr_send.cpp +++ b/src/client/msgr_send.cpp @@ -189,7 +189,11 @@ bool osd_messenger_t::try_send(osd_client_t *cl) } if (ringloop && !use_sync_send_recv) { - io_uring_sqe* sqe = ringloop->get_sqe(); + auto iothread = iothreads.size() ? iothreads[peer_fd % iothreads.size()] : NULL; + io_uring_sqe sqe_local; + ring_data_t data_local; + sqe_local.user_data = (uint64_t)&data_local; + io_uring_sqe* sqe = (iothread ? &sqe_local : ringloop->get_sqe()); if (!sqe) { return false; @@ -200,6 +204,10 @@ bool osd_messenger_t::try_send(osd_client_t *cl) ring_data_t* data = ((ring_data_t*)sqe->user_data); data->callback = [this, cl](ring_data_t *data) { handle_send(data->res, cl); }; my_uring_prep_sendmsg(sqe, peer_fd, &cl->write_msg, 0); + if (iothread) + { + iothread->add_sqe(sqe_local); + } } else { diff --git a/src/osd/osd.cpp b/src/osd/osd.cpp index b8af5f15..75d27359 100644 --- a/src/osd/osd.cpp +++ b/src/osd/osd.cpp @@ -141,6 +141,14 @@ void osd_t::parse_config(bool init) config = msgr.merge_configs(cli_config, file_config, etcd_global_config, etcd_osd_config); if (config.find("log_level") == this->config.end()) config["log_level"] = 1; + if (init) + { + // OSD number + osd_num = config["osd_num"].uint64_value(); + if (!osd_num) + throw std::runtime_error("osd_num is required in the configuration"); + msgr.osd_num = osd_num; + } if (bs) { auto bs_cfg = json_to_bs(config); @@ -150,11 +158,6 @@ void osd_t::parse_config(bool init) msgr.parse_config(config); if (init) { - // OSD number - osd_num = config["osd_num"].uint64_value(); - if (!osd_num) - throw std::runtime_error("osd_num is required in the configuration"); - msgr.osd_num = osd_num; // Vital Blockstore parameters bs_block_size = config["block_size"].uint64_value(); if (!bs_block_size) diff --git a/src/util/ringloop.cpp b/src/util/ringloop.cpp index 410db51b..3c4ad52b 100644 --- a/src/util/ringloop.cpp +++ b/src/util/ringloop.cpp @@ -10,8 +10,9 @@ #include "ringloop.h" -ring_loop_t::ring_loop_t(int qd) +ring_loop_t::ring_loop_t(int qd, bool multithreaded) { + mt = multithreaded; int ret = io_uring_queue_init(qd, &ring, 0); if (ret < 0) { @@ -64,6 +65,25 @@ void ring_loop_t::unregister_consumer(ring_consumer_t *consumer) } } +io_uring_sqe* ring_loop_t::get_sqe() +{ + if (mt) + mu.lock(); + if (free_ring_data_ptr == 0) + { + if (mt) + mu.unlock(); + return NULL; + } + struct io_uring_sqe* sqe = io_uring_get_sqe(&ring); + assert(sqe); + *sqe = { 0 }; + io_uring_sqe_set_data(sqe, ring_datas + free_ring_data[--free_ring_data_ptr]); + if (mt) + mu.unlock(); + return sqe; +} + void ring_loop_t::loop() { if (ring_eventfd >= 0) @@ -79,6 +99,8 @@ void ring_loop_t::loop() struct io_uring_cqe *cqe; while (!io_uring_peek_cqe(&ring, &cqe)) { + if (mt) + mu.lock(); struct ring_data_t *d = (struct ring_data_t*)cqe->user_data; if (d->callback) { @@ -90,12 +112,16 @@ void ring_loop_t::loop() dl.res = cqe->res; dl.callback.swap(d->callback); free_ring_data[free_ring_data_ptr++] = d - ring_datas; + if (mt) + mu.unlock(); dl.callback(&dl); } else { fprintf(stderr, "Warning: empty callback in SQE\n"); free_ring_data[free_ring_data_ptr++] = d - ring_datas; + if (mt) + mu.unlock(); } io_uring_cqe_seen(&ring, cqe); } diff --git a/src/util/ringloop.h b/src/util/ringloop.h index 0f14e2c0..e1945526 100644 --- a/src/util/ringloop.h +++ b/src/util/ringloop.h @@ -14,6 +14,7 @@ #include #include #include +#include #define RINGLOOP_DEFAULT_SIZE 1024 @@ -124,28 +125,21 @@ class ring_loop_t std::vector> immediate_queue, immediate_queue2; std::vector consumers; struct ring_data_t *ring_datas; + std::mutex mu; + bool mt; int *free_ring_data; unsigned free_ring_data_ptr; bool loop_again; struct io_uring ring; int ring_eventfd = -1; public: - ring_loop_t(int qd); + ring_loop_t(int qd, bool multithreaded = false); ~ring_loop_t(); void register_consumer(ring_consumer_t *consumer); void unregister_consumer(ring_consumer_t *consumer); int register_eventfd(); - inline struct io_uring_sqe* get_sqe() - { - if (free_ring_data_ptr == 0) - return NULL; - struct io_uring_sqe* sqe = io_uring_get_sqe(&ring); - assert(sqe); - *sqe = { 0 }; - io_uring_sqe_set_data(sqe, ring_datas + free_ring_data[--free_ring_data_ptr]); - return sqe; - } + io_uring_sqe* get_sqe(); inline void set_immediate(const std::function cb) { immediate_queue.push_back(cb);