Версия для печати темы
Нажмите сюда для просмотра этой темы в оригинальном формате
Форум программистов > C/C++: Общие вопросы > lock-free queues


Автор: J0ker 8.4.2009, 22:16
народ
проверьте пжалста - нигде я тут не напортачил
Код

#include <malloc.h>

template <typename T>
class atomic;

template <typename T>
class atomic<T *>
{
private:
    volatile LONG *_value;
public:
    explicit atomic(const T *value = NULL)
    {
        _value = reinterpret_cast<volatile LONG *>(_aligned_malloc(sizeof(LONG), 32));
        if(_value)
            *_value = *reinterpret_cast<LONG *>(&value);
    }
    ~atomic()
    {
        if(_value)
            _aligned_free((void *)_value);
    }

    atomic(const atomic &a)
    {
        if(_value)
            operator=(a);
    }
    atomic &operator=(const atomic &a)
    {
        if(this != &a)
            ::InterlockedExchange(_value, ::InterlockedExchange(a._value, *a._value));
        return *this;
    }

    operator T*() const
    {
        LONG tmp = ::InterlockedExchange(_value, *_value);
        return *reinterpret_cast<T **>(&tmp);
    }
    atomic &operator=(const T *value)
    {
        ::InterlockedExchange(_value, *reinterpret_cast<LONG *>(&value));
        return *this;
    }

    bool operator==(const atomic &right) const
    {
        return operator T*() == right->operator T*();
    }
    bool operator!=(const atomic &right) const
    {
        return !operator==(right);
    }

    T *operator->()
    {
        return operator T*();
    }

    T *exchange(const T *value)
    {
        LONG tmp = ::InterlockedExchange(_value, *reinterpret_cast<LONG *>(&value));
        return *reinterpret_cast<T **>(&tmp);
    }
};


// multiple producers - single consumer lock-free queue
template <typename T>
class mpsc_queue
{
private:
    struct node
    {
        const T _value;
        atomic<node *> _next;
        explicit node(const T &value): _value(value), _next(NULL) {}
    };
private:
    node *_head;
    atomic<node *> _tail;

public:
    mpsc_queue()
    {
        _head = _tail = new node(T());
    }
    ~mpsc_queue()
    {
        while(_head != NULL)
        {
            node *tmp = _head;
            _head = tmp->_next;
            delete tmp;
        }
    }

    void push(const T &value)
    {
        node *tmp = _tail.exchange(new node(value));
        tmp->_next = _tail;
    }
    bool pop(T &result)
    {
        if(_head->_next == NULL)
            return false;
        node *tmp = _head;
        result = _head->_next->_value;
        _head = _head->_next;
        delete tmp;
        return true;
    }
};

// single producers - single consumer lock-free queue
template <typename T>
class spsc_queue
{
private:
    struct node
    {
        const T _value;
        node *_next;
        explicit node(const T &value): _value(value), _next(NULL) {}
    };
private:
    node *_head;
    atomic<node *> _tail;

public:
    spsc_queue()
    {
        _head = _tail = new node(T());
    }
    ~spsc_queue()
    {
        while(_head != NULL)
        {
            node *tmp = _head;
            _head = tmp->_next;
            delete tmp;
        }
    }

    void push(const T &value)
    {
        _tail->_next = new node(value);
        _tail = _tail->_next;
    }
    bool pop(T &result)
    {
        if(_head == _tail)
            return false;
        node *tmp = _head;
        result = _head->_next->_value;
        _head = _head->_next;
        delete tmp;
        return true;
    }
};

Автор: Lazin 8.4.2009, 22:21
Цитата(J0ker @  8.4.2009,  22:16 Найти цитируемый пост)
    volatile LONG *_value;

в хипе переменные размещать не обязательно smile 
можно использовать ф-ию InterlockedPointerExchange, в которую передавать сразу, указатель на нод

Добавлено через 5 минут и 6 секунд
в классе mpsc_queue
Цитата(J0ker @  8.4.2009,  22:16 Найти цитируемый пост)
    void push(const T &value)
    {
        node *tmp = _tail.exchange(new node(value));
        tmp->_next = _tail;
    }

это race condition, здесь нужно использовать CAS
алгоритм должен быть такой
создаем новый нод tmp, далее пытаемся указатель на него записать в tail, если tail не изменился за время создания нода, то продолжаем работать. если изменился, то заново получаем указатель на tail и опять пытаемся присвоить туда указатель на новый нод, как-то так
а у тебя в коде нигде не используется InterlockedCompareExchange, без этого lock-free queue реализовать невозможно

Автор: J0ker 8.4.2009, 22:27
Цитата(Lazin @ 8.4.2009,  22:21)
Цитата(J0ker @  8.4.2009,  22:16 Найти цитируемый пост)
    volatile LONG *_value;

в хипе переменные размещать не обязательно smile 
можно использовать ф-ию InterlockedPointerExchange, в которую передавать сразу, указатель на нод

спасибо за информацию - потом поправлю
но прежде всего проверьте логику

Добавлено через 7 минут и 7 секунд
Цитата(Lazin @ 8.4.2009,  22:21)
Добавлено @ 22:26
в классе mpsc_queue
Цитата(J0ker @  8.4.2009,  22:16 Найти цитируемый пост)
    void push(const T &value)
    {
        node *tmp = _tail.exchange(new node(value));
        tmp->_next = _tail;
    }

это race condition, здесь нужно использовать CAS
алгоритм должен быть такой
создаем новый нод tmp, далее пытаемся указатель на него записать в tail, если tail не изменился за время создания нода, то продолжаем работать. если изменился, то заново получаем указатель на tail и опять пытаемся присвоить туда указатель на новый нод, как-то так
а у тебя в коде нигде не используется InterlockedCompareExchange, без этого lock-free queue реализовать невозможно

вы не совсем поняли логику
создается овый нод, после чего атомарно берется текущий хвост (помещается в tmp) и указаталь на хвост заменяется новым нодом
если в этот момент произойдет повторный вход в push, то следующий нод будет рисоединен к этому в любом случае

Автор: Lazin 8.4.2009, 22:54
вот как я это вижу
Код

template<class value_type>
struct node_type
{
    value_type value;
    volatile node_type* next;
    node_type(const value_type& v) : value(v), next(0){}
};

template<class value_type>
class mpsc_queue
{
    typedef node_type<value_type> node_t;
    volatile node_t* tail_;
    volatile node_t* head_;

    void enqueue(const value_type& value)
    {
        //create new
        node_t* tmp = new node_t(value);
    loop:
        //сохраняем указатель на элемент
        node_t* old_tail = tail_->next;//вообще, здесь то-же нужно использовать Interlocked ф-ю
        if ((node_t*)InterlockedPointerCompareExchange((LONG**)&tail_->next, (LONG*)tmp, (LONG*)old_tail) != old_tail)
        {
            //другой поток успел добавить элемент в очередь раньше нас
            goto loop;
        }
        //все ок, мы добавили элемент в конец очереди
    }
    ....
};


Цитата(J0ker @  8.4.2009,  22:27 Найти цитируемый пост)
создается овый нод, после чего атомарно берется текущий хвост (помещается в tmp) и указаталь на хвост заменяется новым нодом
если в этот момент произойдет повторный вход в push, то следующий нод будет рисоединен к этому в любом случае

ты прав, должно работать

Добавлено через 4 минуты и 48 секунд
вроде-бы ничего не смущает, должно работать smile 

Автор: Lazin 8.4.2009, 23:20
Цитата(J0ker @  8.4.2009,  22:16 Найти цитируемый пост)
Код

     void push(const T &value)
    {
        node *tmp = _tail.exchange(new node(value));
        tmp->_next = _tail;
    }

мне кажется возможным такой вариант, первый поток входит в push, изменяет tail_ и сохраняет его старое значение в tmp
второй поток входит в push и делает то-же самое, изменяет tail_ и сохраняет в tmp значение которое туда записал другой поток
первый поток просыпается и записывает в tmp->next новое значение tail_, а не то, которое он туда установил
второй поток просыпается и записывает в свой tmp->next (tmp - нод созданый первым потоком, на который никто не указывает в данный момент) указатель на tail_ (тот-же tail_ что и у первого потока)
в результате - утечка памяти

Автор: J0ker 9.4.2009, 00:13
Цитата(Lazin @ 8.4.2009,  23:20)
Цитата(J0ker @  8.4.2009,  22:16 Найти цитируемый пост)
Код

     void push(const T &value)
    {
        node *tmp = _tail.exchange(new node(value));
        tmp->_next = _tail;
    }

мне кажется возможным такой вариант, первый поток входит в push, изменяет tail_ и сохраняет его старое значение в tmp
второй поток входит в push и делает то-же самое, изменяет tail_ и сохраняет в tmp значение которое туда записал другой поток
первый поток просыпается и записывает в tmp->next новое значение tail_, а не то, которое он туда установил
второй поток просыпается и записывает в свой tmp->next (tmp - нод созданый первым потоком, на который никто не указывает в данный момент) указатель на tail_ (тот-же tail_ что и у первого потока)
в результате - утечка памяти

о блин
точно
непрально выразил мысль
должно быть так:
Код

void push(const T &value)
{
    node *new_node = new node(value);
    node *tmp = _tail.exchange(new_node);
    tmp->_next = new_node;
}


Автор: Lazin 9.4.2009, 08:33
Цитата(J0ker @  9.4.2009,  00:13 Найти цитируемый пост)
Код

void push(const T &value)
{
    node *new_node = new node(value);
    node *tmp = _tail.exchange(new_node);
    tmp->_next = new_node;
}


так будет работать, никаких проблем я не вижу smile

Добавлено через 2 минуты и 28 секунд
хотя. можно сделать и mpmc_queue smile 
tbb::concurrent_queue позволяет читать и записывать многим потокам одновременно, может не стоит изобретать свой вид транспрота а использовать threading building blocks? smile 

Автор: J0ker 9.4.2009, 15:49
Цитата(Lazin @  9.4.2009,  08:33 Найти цитируемый пост)
хотя. можно сделать и mpmc_queue 

да, возможно добавлю - пока мне это не актуально

Цитата(Lazin @  9.4.2009,  08:33 Найти цитируемый пост)
tbb::concurrent_queue позволяет читать и записывать многим потокам одновременно, может не стоит изобретать свой вид транспрота а использовать threading building blocks?

это я смотрел
во-первых коммерческая лицензия у них 300 баксов
а во-вторых у них очередь локирующая - они мотивируют это тем, что во-первых такой нагрузки на очередь, что-бы лак сказывался на производительности все равно не бывает, т.к. очереди обычно присутствуют там, где есть IO; а во-вторых у консьюмера все равно будет либо беспробудный пулинг, либо нужен сторонний механизм нотификации, который сам по себе меет всегда локирующую природу
они, конечно, правы, но мне больше нравятся красивые решения, а не красивые отмазки  smile 

Автор: Lazin 9.4.2009, 16:30
Цитата(J0ker @  9.4.2009,  15:49 Найти цитируемый пост)
во-первых коммерческая лицензия у них 300 баксов

в коммерческих проектах можно использовать и бесплатную версию, к тому-же 300$ не особо крупная сумма smile

Добавлено через 1 минуту и 32 секунды
Цитата(J0ker @  9.4.2009,  15:49 Найти цитируемый пост)
они, конечно, правы, но мне больше нравятся красивые решения, а не красивые отмазки

просто там намного больше возможностей (к примеру определение размера очереди), но если эти возможности не нужны, то можно вполне обойтись и "красивым" решением

Powered by Invision Power Board (http://www.invisionboard.com)
© Invision Power Services (http://www.invisionpower.com)