Брокер сообщений на Rust

Всем привет. Написал бессерверный брокер сообщений, кому интересно прошу под кат.

Проект представляет из себя динамическую библиотеку с Си-интерфейсом.

Библиотеку назвал liner (в репе крейтов раста пришлось название сделать пошире liner_broker, поскольку занято имя уже было).

Следующие фичи в наличии сейчас:

Передача данных внутри по TCP, без шифрования.

По плану статьи: посмотрим как реализована библиотека, напишем пример отправки 10к сообщений и сравним общее время отправки-получения с либой ZeroMQ.

Сразу посмотрим как выглядит пример отправки-получения сообщений на расте:

use liner_broker::Liner;

fn  main() {

    let mut client1 = Liner::new("client1", "topic_client1", "localhost:2255", "redis://localhost/");
    let mut client2 = Liner::new("client2", "topic_client2", "localhost:2256", "redis://localhost/");
   
    client1.run(Box::new(|_to: &str, _from: &str, _data: &[u8]|{
        println!("receive_from {}", _from);
    }));
    client2.run(Box::new(|_to: &str, _from: &str, _data: &[u8]|{
        println!("receive_from {}", _from);
    }));

    let array = [0; 100];
    for _ in 0..10{
        client1.send_to("topic_client2", array.as_slice());
        println!("send_to client2");       
    }
}

Клиентов может быть несколько в одном пользовательском процессе. При создании клиента надо задать уникальное имя клиента, название топика, адрес клиента и адрес БД Redis (далее БД). При запуске клиента в работу надо задать колбек-функцию для получения данных, сообщения от всех отправителей будут приходить в нее.

Архитектура библиотеки liner
Архитектура библиотеки liner

Внутри либа состоит из двух крупных блоков: отправитель и получатель сообщений.

Отправитель - это отдельный поток для отправки сообщений клиентов, чтобы не тормозить клиентов и увеличить пропускную способность:

С отправителем закончили. Еще стоит добавить, что сообщения размером более 1мб (зашитый параметр компиляции) сжимаются, используется крейт zstd.

Получатель сообщений - здесь уже два потока, один для ожидания данных в сокетах,
второй - поток обработки, для вызова пользовательских колбеков:

Мемпул собственной разработки (посмотрел что есть, подходящего прямо без переделок не нашел). Сначала без мемпула было, профилировал это все дело, много вызовов было выделения памяти. Работает мемпул таким образом:

Для гарантии доставки сообщений используется БД Redis:

Сейчас посмотрим пример отправки 10к сообщений, каждое размером по 1024байт.

use std::time::SystemTime;
use std::{thread, time};
use std::sync::{ Arc, Mutex};

use liner_broker::Liner;

fn  main() {

    let mut client1 = Liner::new("client1", "topic_client1", "localhost:2255", "redis://localhost/");
    let mut client2 = Liner::new("client2", "topic_client2", "localhost:2256", "redis://localhost/");
   
    client1.clear_stored_messages();
    client2.clear_stored_messages();

    client1.clear_addresses_of_topic();
    client2.clear_addresses_of_topic();

    const MESS_SEND_COUNT: usize = 10000;
    const MESS_SIZE: usize = 1024;
    const SEND_CYCLE_COUNT: usize = 10;

    let mut receive_count: i32 = 0;
    let send_end = Arc::new(Mutex::new(0));
    let _send_end = send_end.clone();

    client1.run(Box::new(|_to: &str, _from: &str,  _data: &[u8]|{
        println!("receive_from {}", _from);
    }));
    client2.run(Box::new(move |_to: &str, _from: &str,  _data: &[u8]|{
        receive_count += 1;    
        if receive_count == MESS_SEND_COUNT as i32{
            receive_count = 0;
            println!("receive_from {} ms", current_time_ms() - *_send_end.lock().unwrap());
        }
    }));

    let array = [0; MESS_SIZE];
    for _ in 0..SEND_CYCLE_COUNT{
        let send_begin = current_time_ms();
        for _ in 0..MESS_SEND_COUNT{
            client1.send_to("topic_client2", array.as_slice());
        }
        let send = current_time_ms();
        println!("send_to {} ms", send - send_begin);  
        *send_end.lock().unwrap() = send;     
    
        thread::sleep(time::Duration::from_millis(1000));
    }
}

fn current_time_ms()->u64{ 
    SystemTime::now()
    .duration_since(SystemTime::UNIX_EPOCH)
    .unwrap()
    .as_millis() as u64
}
alex@ubuntu2004:~/projects/rust/liner/cpp$ cargo build --release
   Compiling liner_broker v1.1.2 (/home/alex/projects/rust/liner)
    Finished `release` profile [optimized] target(s) in 5.62s
alex@ubuntu2004:~/projects/rust/liner/cpp$ cd ../target/release/
alex@ubuntu2004:~/projects/rust/liner/target/release$ ./throughput_10k 
send_to 8 ms
receive_from 8 ms
send_to 5 ms
receive_from 5 ms
send_to 7 ms
receive_from 3 ms
send_to 11 ms
receive_from 3 ms
send_to 6 ms
receive_from 3 ms
send_to 3 ms
receive_from 4 ms
send_to 6 ms
receive_from 4 ms
send_to 12 ms
receive_from 3 ms
send_to 7 ms
receive_from 4 ms
send_to 7 ms
receive_from 3 ms

В итоге в среднем имеем 10мс общего времени.

В плюсовом клиенте медленнее получилось в 2 раза, думаю из-за косвенного вызова колбека. Посмотрим код плюсового примера ради интереса:

#include "liner_broker.h"

#include <ctime>
#include <iostream>
#include <chrono>
#include <thread>

const int MESS_SEND_COUNT = 10000;
const int MESS_SIZE = 1024;
const int SEND_CYCLE_COUNT = 30;

int main(int argc, char* argv[])
{  
    auto client1 = LinerBroker("client1", "topic_client1", "localhost:2255", "redis://localhost/");
    auto client2 = LinerBroker("client2", "topic_client2", "localhost:2256", "redis://localhost/");
 
    int receive_count = 0;
    clock_t send_begin = clock();
    clock_t send_end = clock();

    client1.run([](const std::string& to, const std::string& from, const std::string& data){});
    client2.run([&receive_count, &send_end](const std::string& to, const std::string& from, const std::string& data){
        receive_count += 1;
        if (receive_count == MESS_SEND_COUNT){
            receive_count = 0;
            std::cout << "receive_from " << 1000.0 * (clock() - send_end) / CLOCKS_PER_SEC << " ms" << std::endl;
        }
    });
    
    char data[MESS_SIZE];
    for (int i = 0; i < SEND_CYCLE_COUNT; ++i){
        send_begin = clock();
        for (int j = 0; j < MESS_SEND_COUNT; ++j){
            client1.sendTo("topic_client2", data);
        }
        send_end = clock();
        std::cout << "send_to " << 1000.0 * (send_end - send_begin) / CLOCKS_PER_SEC << " ms" << std::endl;
        std::this_thread::sleep_for(std::chrono::milliseconds(1000));
    }    
}

Теперь запустим аналогичный пример для ZeroMQ

#include <iostream>
#include <zmq_addon.hpp>
#include <ctime>
#include <chrono>
#include <thread>

int main()
{
    zmq::context_t ctx;
    zmq::socket_t sock1(ctx, zmq::socket_type::push);
    zmq::socket_t sock2(ctx, zmq::socket_type::pull);
    sock1.bind("tcp://127.0.0.1:*");
    const std::string last_endpoint =
        sock1.get(zmq::sockopt::last_endpoint);
    std::cout << "Connecting to "
              << last_endpoint << std::endl;
    sock2.connect(last_endpoint);

    std::vector<zmq::const_buffer> send_msgs;
    char mess[1024];
    for (int i = 0; i < 10000; ++i){
        send_msgs.push_back(zmq::str_buffer(mess));
    }
    for (int i = 0; i < 10; ++i){
        auto send_begin = clock();
    
        if (!zmq::send_multipart(sock1, send_msgs))
            return 1;

        std::vector<zmq::message_t> recv_msgs;
        const auto ret = zmq::recv_multipart(
            sock2, std::back_inserter(recv_msgs));
        if (!ret)
            return 1;
        
        auto receive_end = clock();
        std::cout << "send_to " << 1000.0 * (receive_end - send_begin) / CLOCKS_PER_SEC << " ms" << std::endl;
        std::this_thread::sleep_for(std::chrono::milliseconds(1000));
    }
    return 0;
}
alex@ubuntu2004:~/projects/rust/liner$ cd benchmark/compare_with_zeromq/
alex@ubuntu2004:~/projects/rust/liner/benchmark/compare_with_zeromq$ make
g++ -Wall -O2 -std=c++17 -g -Wno-write-strings -o compare_with_zmq compare_with_zmq.cpp -lzmq
alex@ubuntu2004:~/projects/rust/liner/benchmark/compare_with_zeromq$ ./compare_with_zmq 
Connecting to tcp://127.0.0.1:34079
send_to 20.198 ms
send_to 16.504 ms
send_to 11.5 ms
send_to 13.153 ms
send_to 10.964 ms
send_to 10.788 ms
send_to 10.785 ms
send_to 11.119 ms
send_to 11.348 ms
send_to 10.826 ms

Имеем примерно тоже самое. Скорее всего ZeroMQ еще можно как-то подкрутить (сокеты на пайпы заменить, буферами внутренними поиграться и тп) и он быстрее будет думаю.

Ну вот и все пожалуй.

Пока в планах следующие задачки:

Присоединяйтесь к разработке кто хочет-может, или продвижению, если есть интерес к этому делу.

Спасибо.

PS: лицензия MIT, ссылка на гитхаб

@Tyiler
03.03.2025 15:26 UTC
Первоисточник

Комментарии

@PrinceKorwin
03.03.2025 10:37 UTC
+1

гарантия доставки - хотя бы один раз (at-least-once), используется БД Redis для этого

Redis имеет pub/sub. Понимаю, что у вас функций побольше. Но почему не Tarantul? Если правильно понял, то у тарантула есть все, что вы описываете.

@Tyiler
03.03.2025 10:52 UTC
0

Tarantul за место редиса имеете ввиду? Ну не знаю, можно, наверно.

Или вообще не писать ничего, а тарантул брать и пошел? Дык это понятно, много всего готового уже есть, хотелось что-то свое создать.

@aystarik
03.03.2025 12:03 UTC
+1

А чем "БД Редис" не сервер?

@Tyiler
03.03.2025 12:30 UTC
0

Ее опциональной сделаю попозже.

Но без нее возникнут накладные расходы для обеспечения гарантии доставки: надо будет подверждения обратно слать после обработки сообщения (пусть не каждого, но всеравно), на той стороне нужен код, который это будет все слушать и тд

Скорее всего будет так: если без БД, то и без гарантии.

04.03.2025 03:28 UTC
0

Redis тоже не дает особо гарантии. Он может быть in memory only, илм записью на диск, но все равно иметь возможности потери данных.

04.03.2025 06:56 UTC
0

Редис выбрал, потому что схему данных не надо создавать заранее у него, как в обычных БД, и еще потому что он обычно уже есть, то есть много кто им пользуется.

А так да, если не позаботишься сам специально, то потерять данные можно.

04.03.2025 07:15 UTC
+1

Не стоит писать про гарантии доставки/подтверждения если не указываете какие лоя этого должны быть настройки у Redis. Я про это.

@
03.03.2025 13:10 UTC
0
НЛО прилетело и опубликовало эту надпись здесь
@Tyiler
03.03.2025 13:37 UTC
0

В итоге в среднем имеем 10мс общего времени.

10к/10мс отсюда предположим линейно, что будет 1Млн сооб/сек

Ну дык у вас же не млн сообщений в сек. Или я не понял претензии или вопроса.

На бенчмарк я там ссылку дал, все там понятно написано, 50 строчек кода и Makefile тут же рядом.

Возможно у вас сервер мощный с кучей ядер и прочее (если вы считаете, что это очень мало), я тестил на своем рабочем ПК (core 11700).

03.03.2025 15:15 UTC
0

а немного ждет (по умолч 1мс

Про это напишу еще.

Здесь упор сделан на bw (пропуск-ю способность), без задержки этой можно обойтись, обнулить при сборке, если нужно в первую очередь быстродействие.

Без нее часто будет выходить из потока отправки данных чуть раньше времени, чем новая порция сообщений для отправки придет. Все измерено было, с задержкой этой и без, то есть ее не просто так добавил.