Страницы

Поиск по вопросам

Показаны сообщения с ярлыком очередь. Показать все сообщения
Показаны сообщения с ярлыком очередь. Показать все сообщения

понедельник, 9 марта 2020 г.

Как убрать из математического выражения лишние скобки?

#cpp #алгоритм #очередь


Попалась такая задача: «при помощи очереди изъять лишние скобки из арифметического
выражения». Изначально она показалась очень простой, но на практике вызвала много затруднений.

Добавление @mymedia:

Лишними скобками считаются те, которые можно убрать, сохранив смысл выражения. Например,
(5*7)+3 → 5*7+3, ((a+b)) → a+b.
    


Ответы

Ответ 1



Задача действительно не самая простая и через стек ее не решить. Через стек решается задача определения правильности/неправильности скобочной последовательности. Лишние скобки так не убрать. ВОТ можно почитать про обратную польскую нотацию. Раскладываете выражение, собираете обратно и где надо ставите скобки. Тогда лишних уже не будет. P.S. очередь/стек там используется для сохранения операций.

понедельник, 24 февраля 2020 г.

Будет ли равномерным распределение заданий в очереди?

#python #многопоточность #очередь


Нужно что-то вроде RR.

Будет ли равномерным распределение заданий в очереди Queue.Queue() по тредам? Смущает
то что задания появляться будут медленней чем выполнение команд по ним.

Запускаю несколько тредов (с разным кодом под разные устройства) и тред который принимает
задания по http. Задание попадает в очередь, а до  этого воркеры висят заблокированном
queue.get().

Я использовал в другом проекте multiprocessing.dummy.Pool - там все красиво и равномерно,
но сейчас треды имеют разный код внутри и одной функцией не обойдешься.

Задача разгрузить исполнительные устройсва, а не процессор.
    


Ответы

Ответ 1



Сперва почему-то хотелось ответить, что распределение неравномерное - то есть часть потоков будет простаивать, а тот, что был создан первее всех будет отдуваться за всех. Видимо такой поспешный вывод из-за наивного представления потоков в воображении. Но вот эксперимент (py3) - задания появляются в два раза медленнее, чем обрабатываются: import queue import threading import time import random counter = {} lock = threading.Lock() def worker(): while True: try: sleep = q.get() if sleep == "DIE": q.task_done() return print(threading.current_thread().name) with lock: counter[threading.current_thread().name] = counter.get(threading.current_thread().name, 0) + 1 time.sleep(sleep) q.task_done() except queue.Empty: pass q = queue.Queue() pool = [threading.Thread(target=worker, name="Worker " + str(i)) for i in range(4)] [t.start() for t in pool] for i in range(100): rand_sleep = random.random() / 8 q.put(rand_sleep) time.sleep(rand_sleep * 2) for i in range(4): q.put("DIE") q.join() print("Report: ", counter) Вот что скрипт выводит: output: ('Report: ', {'Worker 2': 25, 'Worker 3': 25, 'Worker 0': 25, 'Worker 1': 25}) Как видно из результатов - для моего кода распределение равномерное. Если задания поступают без задержек - то распределение неравномерное и зависит от того, как долго обрабатывает задание каждый поток. Дополнение: данный скрипт был протестирован в Windows и Ubuntu, использовался Python3.4 и Python2.7. Во всех четырех случаях результаты одинаковы.

вторник, 28 января 2020 г.

Что такое стеки и очереди? [закрыт]

#cpp #очередь #стек


        
             
                
                    
                        
                            Закрыт. Данный вопрос необходимо конкретизировать. Ответы
на него в данный момент не принимаются.
                            
                        
                    
                
                            
                                
                
                        
                            
                        
                    
                        
                            Хотите улучшить этот вопрос? Переформулируйте вопрос,
чтобы он был сосредоточен только на одной проблеме, отредактировав его.
                        
                        Закрыт 9 месяцев назад.
                                                                                
           
                
        
Учусь на втором курсе, программирование, C++. Расскажите, пожалуйста, о стеках. Как
пишется программа?
    


Ответы

Ответ 1



Стек - структура данных с доступом к элементам по принципу LIFO (Last In First Out - Последний пришел - первый вышел). Данные добавляются в начало (конец, кому как удобно), оттуда же и извлекаются. Для реализации данной структуры достаточно иметь лишь две функции и указатель на верхушку: push(item); //добавляет данные в стек pop(); //извлекает последний элемент Очередь - структура данных с доступом к элементам по принципу FIFO (First In First Out - Первый пришел - Первый вышел). Данные добавляются в конец, а извлекаются из начала. Для быстрого добавления и извлечения данных понадобится два указателя, один на начало, очереди, второй на ее конец. Для работы с очередью используются функции: enqueue(item); //добавляет новый элемент в очередь dequeue(); //извлекает элемент из очереди Для работы с указателями в C++ вам необходимо ознакомиться, например, с данной статьей (одна из первых в поиске). Затем, вам необходимо создать элемент данных вашего стека (очереди). Например: struct data { int p; int c; } После этого обзавестись ячейкой стека. Например: struct element { data element_data; element *next; //указатель на следующий элемент } Указатель на следующий элемент необходим дабы наши данные в стеке (очереди) имели некоторую последовательность и были как-то связаны между собой. После этого пишете необходимые функции, обзаводитесь указателями на начало (и конец, при необходимости). Например для push(): //somewhere in the code element *top; //инициализируется в NULL //... void push(data newData) { //создаем нашу структуру element и указатель на нее element *newElement = new element(); //element_data = newData element->element_data = newData; //next = top - здесь мы и задаем связку между двумя элементами стека element->next = top; //top = указателю на нашу структуру top = newElement; } Все остальное по аналогии. На C++ давненько я не кодил (проверить негде), поэтому прошу заранее прощения всех гуру, поправьте если где напортачил.

Ответ 2



Я бы описал стэк как упаковку с коллекцией монеток, уложеных в стопочку. С 1 до 10. Что бы получить доступ к 5, нужно снять 10, потом 9, 8, 7, 6 и только потом получим 5-ю. Ну это лично моё ИМХО, когда то мне было проще запомнить так. В остальном остаются только технические вопросы, и я не могу не согласится с Dex'ом...

воскресенье, 26 января 2020 г.

C# Файлы, сеть, многопоточность

#c_sharp #многопоточность #сеть #fileupload #очередь


Оговорюсь сразу, использую .NET 4.0. Это в случае будущих рекомендаций в пользу async-await =)

Изначальная постановка задачи - имеется несколько однотипных приложений, функция
которых одна - записывать файлы определенного формата заранее известной длины (например,
10 минут приложение пишет файл, постоянно его увеличивая, по прошествии этого времени
оно закрывает файл и открывает новый). По завершению записи очередного файла оно отправляет
моему приложению TCP сообщение (в качестве TCP сокета использую класс Ipdaemon от ipWorks
от nSoftware) с именем файла, которое оно только что закончила писать. Размер этого
файла может составлять более полутора гигабайт. Задача моего приложения - раз в 10-15
секунд переписывать очередную порцию того, что накатало то приложение, на сетевой диск
и по получению TCP сообщения о завершении - удостовериться в том, что файлы совпадают
по содержимому.

Я написал класс для решения этой проблемы и вот его ключевой метод для перекидывания
(в среднем, его исполнение занимает от 90мс до 8000мс - StopWatch, но теоретически
может и больше):

private void ArrayCopyFile()
    {
        //http://stackoverflow.com/questions/1246899/file-copy-vs-manual-filestream-write-for-copying-file
        //http://stackoverflow.com/questions/995320/file-writeallbytes-causes-error-insufficient-system-resources-exist-to-complete#995320
        const int bufferSize = 1024*1024;
        var headerBytes = new byte[_headerManager.HeaderSize()];
        FileStream src = new FileStream(SrcFile, FileMode.Open, FileAccess.Read,
FileShare.ReadWrite);
        FileStream dst = new FileStream(DstFile, FileMode.OpenOrCreate, FileAccess.Write,
FileShare.ReadWrite);
        try
        {
            // выясним, увеличился ли исходный файл
            if (src.Length > _currPos)
            {
                // чтение заголовка файла
                src.Read(headerBytes, 0, headerBytes.Length);
            }
            else 
            {
                Trace("Исходный файл '{0}' не изменился со времени последней синхронизации",
                    SrcFile);
                return;
            }
            if (!CheckDstDirectory())
            {                    
                return;
            }
            if (_currPos > 0)
            {
                // записываем заголовок
                dst.Write(headerBytes, 0, headerBytes.Length);
                // резервируем место
                dst.SetLength(src.Length);
                // устанавливаем позицию на такую же, как у исходного файла ДО считывания
очередной порции байт
                dst.Seek(_currPos, SeekOrigin.Begin);
            }
            src.Position = _currPos;
            int bytesRead;
            byte[] bytes = new byte[bufferSize];
            while ((bytesRead = src.Read(bytes, 0, bufferSize)) > 0)
            {
                dst.Write(bytes, 0, bytesRead);
                _currPos += bytesRead;
            }
        }
        finally
        {
            src.Dispose();
            dst.Dispose();
        }
    }


Для каждого объекта этого класса создается свой объект System.Threading.Timer с интервалом
10 секунд.

Для проверки идентичности использую такой метод:

public static bool CompareFiles(string filePath1, string filePath2)
    {
        long fileLength;
        if((fileLength = new FileInfo(filePath1).Length) != new FileInfo(filePath2).Length)
            return false;
        bool filesAreEquals = true;
        const int size = 1024*1024; //0x1000000;
        int countIteration = (int)Math.Ceiling(fileLength / (double)size);
        Parallel.For(0, countIteration, x =>
        {
            if(!filesAreEquals) return;
            var start = x * size;
            if (start >= fileLength) return;
            int realSize = (int) (x == countIteration - 1 ? fileLength - start : size);
            using (FileStream file = File.OpenRead(filePath1))
            using (FileStream file2 = File.OpenRead(filePath2))
            {
                var buffer = new byte[realSize];
                var buffer2 = new byte[realSize];
                file.Position = start;
                file2.Position = start;
                int count = file.Read(buffer, 0, realSize);
                file2.Read(buffer2, 0, realSize);
                for (int i = 0; i < count; i++)
                    if (buffer[i] != buffer2[i])
                    {
                        filesAreEquals = false;
                        return;
                    }
            }
        });
        return filesAreEquals;
    }


После приемки TCP сообщения оно помещается в очередь:

class QueueTasks : IDisposable
{
    private readonly Queue _tasks = new Queue();
    readonly object _syncObj = new object();
    private readonly AutoResetEvent _autoResetEvent = new AutoResetEvent(false);
    private readonly ManualResetEvent _exitEvent = new ManualResetEvent(false);
    private bool _isRunning;

    public bool IsRunning
    {
        get { return _isRunning; }
        private set
        {
            _isRunning = value;
            if(!_isRunning) OnQueueTasksStopped();
        }
    }

    public delegate void QueueTasksStoppedEventHandler(object sender);
    public event QueueTasksStoppedEventHandler QueueTasksStopped;
    protected virtual void OnQueueTasksStopped()
    {
        var handler = QueueTasksStopped;
        if (handler != null) handler(this);
    }

    public delegate void TaskReceivedEventHandler(object sender, T eventArgs);
    public event TaskReceivedEventHandler TaskReceived;
    protected virtual void OnTaskReceived(T eventArgs)
    {
        var handler = TaskReceived;
        if (handler != null) handler(this, eventArgs);
    }

    public void Stop()
    {
        _exitEvent.Set();
        _autoResetEvent.Set();
    }

    public void Start()
    {
        if(IsRunning) throw new InvalidOperationException("Очередь задач уже работает");
        IsRunning = true;
        _autoResetEvent.WaitOne();
        _exitEvent.WaitOne();
        new Thread(Work) { Name = "Очередь задач" }.Start();
    }

    public void Add(T data)
    {
        SafeAdd(data);
    }

    void SafeAdd(T data)
    {
        lock (_syncObj) _tasks.Enqueue(data);
        _autoResetEvent.Set();
    }

    T SafeRemove()
    {
        lock (_syncObj)
            if (_tasks.Count > 0)
                return _tasks.Dequeue();
        return default(T);
    }

    void Work()
    {
        while (!_exitEvent.WaitOne(0, false))
        {
            _autoResetEvent.WaitOne();
            while (_tasks.Count != 0)
                OnTaskReceived(SafeRemove());
        }
        IsRunning = true;
    }

    public void Dispose()
    {
        Stop();
    }
}


На событии из этой очереди я анализирую список текущих обрабатываемых файлов, если
нахожу объект класса с именем файла, который совпадает с сообщением, я останавливаю
внутренний таймер для него, проверяю, увеличился ли исходный файл, и если да, докопирую
оставшееся.

Теперь, после всей этой возможно неинтересной лабуды, возможно уже решавшейся и не
раз, задаю вопрос:

Если файлов больше 100, то в очереди начинают копиться сообщения о завершении/начале
записи файлов. Серверная машина, на которой этот софт запускается, имеет 8 ядер чистого
интела и 32 гб оперативной памяти. Пиковая загрузка процессора не выше 80%, оперативная
память - занято около 2-3 гигабайт. Загрузка на сетевом адаптере - не более 55% от
гигабитной сети. Сетевой диск - SSD. В каком месте мое приложение, перекидывающее файлы,
можно и нужно оптимизировать?
    


Ответы

Ответ 1



Parallel.For(0, countIteration, x => Параллельное чтение файлов - это почти всегда плохая идея. Просто читай их в одном потоке последовательно достаточно крупными кусками. И нет смысла каждый раз создавать новый буферный массив - создай по одному на файл и используй. int count = file.Read(buffer, 0, realSize); file2.Read(buffer2, 0, realSize); for (int i = 0; i < count; i++) А это вообще на баг похоже. Stream не обязан прочитать ровно столько байтов, сколько ты ему передал, поэтому Не факт, что файлы одинаковы, если твоя функция так говорит Например, могли сравниться только префиксы блоков. Не факт, что файлы различны, если твоя функция так говорит Например, из первого файла было прочитано больше, чем из второго.

воскресенье, 12 января 2020 г.

Выбор альтернативного сервиса при отказе основного.

#c_sharp #wcf #очередь #cluster


Реализуем кластерное решение для wcf сервисов. Подскажите есть ли какое-нибудь проверенное
решение для проверки доступности другого сервиса и выбора альтернативного, при отказе
основного. 
Например. 
На wcf сервис приходят сообщения от клиента и он их ставит в очередь. Очередь реализовала
на отдельном сервере и отдельным сервисом. Если этот сервис с очередью выйдет из строя,
нужно отправлять сообщения на резервный сервис, определенный заранее. Есть ли какие-то
красивые решения для реализации в автоматическом режиме по средствам wcf или проверять
в коде доступность одного сервиса и так далее?
    


Ответы

Ответ 1



Про штатные средства для подобных проверок в WCF я не слышал (точнее, слышал о WS-Discovery, начиная с .NET 4.0, но не пробовал и не уверен, что оно вообще подходит). Поэтому расскажу про велосипеды. Нужно быстро узнавать о падении стороннего сервиса Суть в периодическом "пинге" сервиса. Заводите таймер/отдельный поток, который проверяет доступность сервиса. Если сервис недоступен, имеет смысл сделать дополнительные 2-3 попытки с возрастающим интервалом между ними (как и в случае со всеми остальными методами), поскольку могут быть кратковременные перебои в сети. Если в итоге сервис не отвечает, производите операцию переключения (меняете адрес сервиса на дополнительный и делаете все остальные вещи, которые могут быть с этим связаны). Повторные попытки, само собой, нужно делать при получении определенных исключений (например, CommunicationException). Контроля над сторонним сервисом нет Если сервис не предлагает специально предназначенной для проверки статуса операции (типа Ping(), Heartbeat(), CheckStatus() и т.д.), то можно запрашивать его метаданные: bool isServiceUp = true; try { string address = "http://server/Service.svc?wsdl"; var mexClient = new MetadataExchangeClient( new Uri(address), MetadataExchangeClientMode.HttpGet); mexClient.GetMetadata(); } catch (Exception e) { // если сервис недоступен, получим исключение isServiceUp = false; } Контроль над сторонним сервисом есть В сервис добавляется "проверочный" метод (Ping()/Heartbeat()/CheckStatus()). В самом простом варианте этот метод пуст, но он также может возвращать и данные о своем состоянии (особенно если внутри себя он использует несколько разных систем - БД, другой сервис и т.д.). О падении сервиса достаточно узнавать в момент его вызова В этом случае все несколько проще, потому что вам не нужно специально проверять сервис. Если при очередном вызове сервиса он не ответил, пробуете еще 2-3 раза. Если после этого сервис не ответил, производите операцию переключения и снова пробуете сделать вызов. Какой бы вариант вы не выбрали, вам понадобятся следующие блоки: код переключения сервисов (как минимум подмена одного адреса на другой) обобщенный код повтора операций, чтобы любой метод можно было вызвать как InvokeWithRetry(() => SomeMethod()) обобщенный код, который соединяет эти два блока вместе и вызывает переключение сервисов в случае необходимости

вторник, 7 января 2020 г.

Thread-safe очередь, и все-все-все

#очередь #c #многопоточность #cpp


Опять у меня идиотский вопрос (простите, но впадаю в старческий маразм, "бабушка
ничего не помнит").
Дано: поток данных поступает в очередь (com-порт). Другой поток разгребает эту очередь,
складывает разобраные данные в базу. Всё в рамках одного приложения.
Вопрос: что юзать в качестве очереди на C (возможно, C++)    


Ответы

Ответ 1



Видимо что-то в таком духе (увидел знакомую тему и не удержался). #include #include #include #include #include struct qitem { struct qitem *next; int len; char data[1]; // реально здесь будет len+1 байт данных }; struct queue { struct qitem *head, *tail; // head == tail == 0 очередь пуста pthread_mutex_t lock; // мьютекс для всех манипуляций с очередью pthread_cond_t cond; // этот cond "сигналим" когда очередь стала НЕ ПУСТОЙ pthread_t th; // tid обработчика }; void inqueue (const char *str, int len, struct queue *q) { struct qitem *p = (typeof(p))malloc(sizeof(*p) + len); strcpy(p->data, str); p->len = len; p->next = 0; pthread_mutex_lock(&q->lock); if (q->tail) q->tail->next = p; else { q->head = p; pthread_cond_broadcast(&q->cond); // теперь очередь не пуста, сигнализируем } q->tail = p; pthread_mutex_unlock(&q->lock); } void processit (struct qitem *p) { int t = atoi(p->data); if (t > 0) { printf ("Sleep %d\n", t); sleep(t); } } // обработчик очереди в отдельном потоке void * consumer (void *arg) { struct queue *q = (typeof(q))arg; for (;;) { pthread_mutex_lock(&q->lock); if (!q->head) { // очередь пуста, делать нечего, ждем... pthread_cond_wait(&q->cond, &q->lock); //pthread_mutex_unlock(&q->lock); // это для нескольких обработчиков очереди //continue; // в нашем случае не нужно // для нескольких "concurrent" обработчиков нужна последовательность // lock/check/wait/unlock/continue/lock/check... } struct qitem *p = q->head; q->head = q->head->next; q->nitems--; if (!q->head) q->tail = q->head; pthread_mutex_unlock(&q->lock); printf ("consume: %s", p->data); if (strcmp(p->data, "STOP") == 0) break; processit(p); free(p); } return 0; } // сделаем пустую очередь и запустим ее обработчик в новом потоке struct queue * run_consumer() { struct queue *q = (typeof(q))malloc(sizeof(*q)); q->head = q->tail = 0; q->nitems = 0; pthread_mutex_init(&q->lock, 0); pthread_cond_init(&q->cond, 0); pthread_create(&q->th, 0, consumer, (void *)q); return q; } int main () { void *res = 0; char *in = NULL; size_t sz; int l; struct queue *q = run_consumer(); while ((l = getline(&in, &sz, stdin)) > 0) { inqueue(in, l, q); in = NULL; } inqueue("STOP", 5, q); if (pthread_join(q->th, &res)) perror("join"); return (long)res; } Для тестирования просто читаем строки с клавиатуры, и ставим их в очередь. Обрааботчик их печатает. Если в начале строки число, то обработчик делает sleep, позволяя написать в очередь несколько строк.

Ответ 2



Я бы использовал просто queue. Поскольку у вас доступ к очереди из нескольких потоков, вам нужно синхронизировать обращения к очереди. Для синхронизации подошёл бы std::mutex, если он доступен на вашем компиляторе. В любом случае, какая-то синхронизация вам всё равно нужна, т.к. в её отсутствие у разных потоков может быть разное представление о содержимом памяти (например потому, что в многопроцессорной системе каждый процессор сбрасывает кэш независимо).

Ответ 3



Используйте boost::lockfree::queue

пятница, 27 декабря 2019 г.

Сервер очередей под задачу

#клиент_сервер #очередь #rabbitmq #backgroundworker #gearman


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

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

Язык программирования, в принципе, не важен, главное чтоб многопоточность поддерживал
(многоядерную).
    


Ответы

Ответ 1



RabbitMQ позволит тебе подключить к одной очереди несколько consumer-ов (см. http://www.rabbitmq.com/tutorials/tutorial-two-python.html Fair Dispatch). По поводу динамических консьюмеров - тут твое приложение само должно решать когда их запускать и как их выключать. Можешь запустить кол-во консьюмеров равное кол-ву ядер на машине и потом собирать данные. А твоя задача напоминает реализацию MapReduce. Может тебе имеет смысл взглянуть на Hadoop? Если все же хочешь написать сам, то по поводу языка можно выбрать Erlang/OTP (на нем написан RabbitMQ). Работа с потоками там построена очень хорошо и объединить несколько машин в кластер, с шарингом потоков, не составит труда. Твое приложение не будет думать о том, где и что выполняется, а просто ждать завершения всех воркеров и вернет результат. Но сломать мозг придется, прежде чем написать функциональный код :)

Ответ 2



Из того, что уже советовали - RabbitMQ. Если нагрузки маленькие и вам необязательна гарантируемость доставки, бурите его. В случае с Rabbit смиритесь с тем, что сообщения могут потеряться, смиритесь с тем, что при больших нагрузках кластер начнёт отказывать, смиритесь с тем, что при разделении сети будут проблемы. Можно ещё взять Apache Kafka. Там с этим лучше, но его настраивать очень сложно, чтоб он хорошо работал. Для вашей же задачи подойдёт Hadoop или Spark. Если хоститесь на амазоне, то с их лямбдами вроде как подобное можно реализовать.

воскресенье, 22 декабря 2019 г.

BlockingCollection - как не блокировать поток

#c_sharp #многопоточность #коллекции #очередь


Применение BlockingCollection, используя подход, когда элементы вытаскиваются из
очереди(например ConcurrentQueue), используя метод Take в цикле -  всегда блокирует
поток. Очевидно, что процессорное время не занимается, однако поток все же занят и
не может использоваться для выполнения других задач. Какая есть альтернатива, когда
нужно последовательно вычитывать элементы из очереди и при этом не блокировать поток?
Конечно, можно сделать велосипед, накрутить событий или чего-нибудь еще, но хотелось
бы понять, нет ли каких-либо стандартных способов это сделать, кроме как использовать
BlockingCollection и метод Take.

P.S. Есть метод TryTake, но я не могу найти решение с его использованием, эквивалентное
использованию Take и при этом неблокирующее поток.
    


Ответы

Ответ 1



Вам на самом деле нужен класс BufferBlock из библиотеки Dataflow (nuget-пакет Microsoft.Tpl.Dataflow). Этот класс заменяет собой BlockingCollection, и позволяет асинхронный доступ: await queue.ReceiveAsync() Таким образом, поток не будет заблокирован. Но у вас получится async-интерфейс. Больше примеров с работающим кодом есть в этом ответе. Ещё одним вариантом является async-обёртка над IProducerCosumerCollection из AsyncEx Стивена Клири: https://github.com/StephenCleary/AsyncEx/wiki/AsyncCollection

вторник, 17 декабря 2019 г.

Очередь задач в PHP

#очередь #php #yii


По мере создания своего проекта столкнулся с проблемой - есть некоторые действия
пользователей, которые могут длиться до нескольких минут. Ясно понятно, что заставлять
пользователя ждать, пока выполняются такие длинные запросы к серверу непозволительно!
И тут пришел к выводу, что нужно организовать очередь задач, чтобы пользователь нажал
ссылку, на сервере сформировалась задача, а пользователю лишь только отображался процесс
выполнения задачи.

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

Слышал о RabbitMQ, Apache Message Queue, но мне кажется, что они слишком избыточны
для моего случая. Мне нужно-то, делать некоторые проверки на сервере (тут как раз пригодился
бы способ выполнения PHP кода, и в зависимости от результата продолжалось бы выполнение
задач, или нет), и манипулировать учетными файлами, и просто "тяжелыми" файлами.

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

Проект создается на базе фреймворка Yii.
    


Ответы

Ответ 1



Воркер на пхп и запускать его по крону. Не?

Ответ 2



ПО мере создания своего проекта столкнулся с проблемой такого плана - есть некоторые действия пользователей, которые могут длиться до нескольких минут. Это ужасно! Это следствие плохой структуры БД и скриптов. Что такого может выполняться, да к тому же у каждого пользователя по минуте по две? Ладно обработка нескольких таблиц с математическими действиями по подсчете скидок у 16К клиентов и 25К заказов. Тут да у меня на серваке до 8 минут считаются, но раз в день и для всех сразу, а не для каждого по отдельности. Я даже набросал табличку в БД, и все бы хорошо, но решение получилось не универсальным. Созданная очередь задач умеет выполнять только шелл скрипты, а мне бы хотелось научить ее выполнять еще логику на php, и чтобы результаты обоих случаев складывались в БД. В Yii есть консольные приложения(и тут), юзай их. Делай там и логику и операции с данными. Не нужно изобретать с данном случае велосипед.

вторник, 10 декабря 2019 г.

Как на C# выполнить метод в “очереди” раз в секунду, а не чаще?

#c_sharp #очередь


Мне нужно, как только происходит определённое событие "Данные обновились"- запускать
метод X -"Разослать обновлённые данные клиентам".
Но если событие вызывается, скажем 50 раз в секунду- нет необходимости столько же
раз сразу же запускать метод X, создавая нагрузку.
Можно ли как-то поставить событие в очередь, и если очередь есть, запустить метод
X- то выполнить только самое последнее событие? 



Вариант 1: Есть 50 клиентов онлайн. Подключается 51-й. Я рассылаю 50+1 клиенту список
клиентов онлайн. Подкл 52-й - делаем то же самое.
Вариант 2: Есть 50 клиентов онлайн. Подключается 51-й. Отсылаю так-же всё. Но тут
подключается пара десятков человек и я должен так же всё рассылать.

Я хочу, если только что была рассылка- не отправлять сразу, а ждать, скажем, секунду
и потом, если был запрос разослать данные- то сделать рассылку таблицы на текущий момент
времени. Таким образом все клиенты получили данные- пусть не моментально, а на секунду
позже, зато я не слал очень много раз.
Ну как-то так.
Как такое реализовать?
    


Ответы

Ответ 1



Ну, вы можете завести таймер и добавить в него немного логики. Вот пример: class Program { public static event EventHandler NewClientArrived; static void X() => Console.WriteLine("\n\tSending notifications"); static void Main(string[] args) { var timer = new System.Timers.Timer(3000) { AutoReset = false }; timer.Elapsed += (sender, timerargs) => X(); NewClientArrived += (sender, eventargs) => { if (!timer.Enabled) timer.Start(); }; // запускаем симуляцию прихода клиентов while (Console.ReadKey().Key != ConsoleKey.Escape) NewClientArrived(null, EventArgs.Empty); } } Если вы используете RX Extensions, можно попробовать так: static void Main(string[] args) { using (Observable.FromEventPattern(typeof(Program), nameof(NewClientArrived)) .Sample(TimeSpan.FromSeconds(1)) .Subscribe(_ => X())) { // запускаем симуляцию прихода клиентов while (Console.ReadKey().Key != ConsoleKey.Escape) NewClientArrived(null, EventArgs.Empty); } }

Ответ 2



Как вариант, можно хранить дату последнего срабатывания события. Когда возбуждается событие, проверяй последний дату выполнения. Если ты видишь, что с момента последнего срабатывания события не прошло и секунды, делай return.

суббота, 30 ноября 2019 г.

Ревью класса для очереди команд

#c_sharp #многопоточность #очередь #очередь_задач #инспекция_кода


Мне потребовалась очередь команд, где множество потоков может добавлять команды на
выполнение и один поток по очереди их выполняет. Тк я не знаю какие-либо стандартные
реализации подобного, то пришлось делать самому. Можете оценить мой класс для этого,
тк у меня сомнения на счет его качества, хотя он работает. И если кто-либо даст ссылку
на реализации подобного, то тоже будет приятно.    

public class ConcurrentCommandsQueue : IDisposable
{
    private readonly object _lockForSyncExecuting;
    private readonly Queue _commandQueue;

    private readonly Action _executeCommand;
    private readonly Action _onStartCommandsExecution;
    private readonly Action _onEndCommandsExecution;

    private volatile bool _isDisposed;
    private volatile bool _isExecuting;
    private volatile Task _taskForAction;

    public bool IsExecuting
    {
        get { return _isExecuting; }
    }

    public ConcurrentCommandsQueue(Action executeCommand) : this(executeCommand,
null, null) { }
    public ConcurrentCommandsQueue(Action executeCommand, 
        Action onStartCommandsExecution, Action onEndCommandsExecution)
    {
        _lockForSyncExecuting = new object();
        _commandQueue = new Queue(10);

        _isDisposed = false;
        _isExecuting = false;
        _executeCommand = executeCommand;
        _onStartCommandsExecution = onStartCommandsExecution;
        _onEndCommandsExecution = onEndCommandsExecution;
        _taskForAction = new TaskCompletionSource().Task;
    }

    private void ExecuteCommands()
    {
        if (_onStartCommandsExecution != null)
            _onStartCommandsExecution();
        T command = default(T);
        while (true)
        {
            lock (_lockForSyncExecuting)
            {
                if (_commandQueue.Count == 0)
                {
                    _isExecuting = false;
                    break;
                }
                command = _commandQueue.Dequeue();
            }
            _executeCommand(command);
        }
        if (_onEndCommandsExecution != null)
            _onEndCommandsExecution(command);
    }

    public void AddCommand(T command)
    {
        if (_isDisposed)
            throw new ObjectDisposedException(GetType().ToString());

        lock (_lockForSyncExecuting)
        {
            if (_isDisposed)
                return;
            _commandQueue.Enqueue(command);
            if (!_isExecuting)
            {
                _isExecuting = true;
                _taskForAction = Task.Run(new Action(ExecuteCommands));
            }
        }
    }
    public void ClearCommandQueue()
    {
        if (_isDisposed)
            throw new ObjectDisposedException(GetType().ToString());

        lock (_lockForSyncExecuting)
            _commandQueue.Clear();
    }
    public void Dispose()
    {
        if (_isDisposed)
            return;

        _isDisposed = true;
        lock (_lockForSyncExecuting)
            _commandQueue.Clear();
        if (_taskForAction != null)
            _taskForAction.Wait();
    }
}

    


Ответы

Ответ 1



У вас получился код в стиле .NET 1.0, который слегка эволюционировал с приходом лямбд (Delegate механически заменён на Action) и задач (new Thread().Start() механически заменён на Task.Run()), но по сути совершенно не изменился: вы TPL с async/await толком не используете, а существование concurrent коллекций вообще упустили из виду. В .NET 4.5 это всё ненужно в принципе. Если вам нужно выполнение задач строго в одном потоке, то можно просто кидать задачи в однопоточный планировщик задач (task scheduler). Собственно, на этой первой строчке весь код и заканчивается. Если нужно выполнение кода перед всеми операциями и после всех операций, то код перед и после этой строчки и пишется. Отмечу, что при этой однострочной реализации вы ещё имеете бонусом возможность использования токена отмены (cancellation token) и прочих радостей жизни. Рассмотрим пример. Допустим, надо выполнить две очереди задач: в каждой очереди задачи выполняются последовательно, но очереди выполняются параллельно. Создаём два планировщика задач с ограничением параллелизма, затем создаём задачи с указанием нужного нам планировщика. class Program { static readonly Random _rnd = new Random(); static readonly LimitedConcurrencyTaskScheduler _schedulerFoo = new LimitedConcurrencyTaskScheduler(1); static readonly LimitedConcurrencyTaskScheduler _schedulerBar = new LimitedConcurrencyTaskScheduler(1); static void Main () => new Program().Run().Wait(); async Task Run () { Task queueFoo = RunQueue("Foo", _schedulerFoo, Enumerable.Range(0, 3).Select(i => (Action)(() => Foo("Foo")))); Task queueBar = RunQueue("Bar", _schedulerBar, Enumerable.Range(0, 3).Select(i => (Action)(() => Foo("Bar")))); await Task.WhenAll(queueFoo, queueBar); Console.WriteLine("Done!"); Console.ReadKey(); } async Task RunQueue (string name, TaskScheduler scheduler, IEnumerable commands) { Console.WriteLine($"{name}: Start"); await Task.WhenAll(commands.Select(c => RunTask(c, scheduler))); Console.WriteLine($"{name}: Finish"); } async Task RunTask (Action task, TaskScheduler scheduler) { await Task.Factory.StartNew(task, CancellationToken.None, TaskCreationOptions.None, scheduler); } void Foo (string name) { int timeout = _rnd.Next(200); Console.WriteLine($"{name}: Start {timeout}"); Thread.Sleep(timeout); Console.WriteLine($"{name}: Finish {timeout}"); } } Пример вывода: Foo: Start Bar: Start Foo: Start 165 Bar: Start 50 Bar: Finish 50 Bar: Start 39 Bar: Finish 39 Bar: Start 115 Foo: Finish 165 Foo: Start 116 Bar: Finish 115 Bar: Finish Foo: Finish 116 Foo: Start 0 Foo: Finish 0 Foo: Finish Done! Хотя в этом примере все задачи и кидаются одновременно, задачи можно добавлять в любой момент. Правда тогда смысл "начала" и "конца" несколько теряется. Если вы объясните, что вы собираетесь с ними делать, тогда можно будет добавить. Здесь я воспользовался планировщиком LimitedConcurrencyTaskScheduler из примеров на MSDN: Samples for Parallel Programming with the .NET Framework (статья с описанием: ParallelExtensionsExtras Tour - #7 - Additional TaskSchedulers). /// /// Provides a task scheduler that ensures a maximum concurrency level while running on top of the ThreadPool. /// Source: http://code.msdn.microsoft.com/ParExtSamples /// Documentation: http://blogs.msdn.com/b/pfxteam/archive/2010/04/09/9990424.aspx /// License: MS-LPL /// public class LimitedConcurrencyTaskScheduler : TaskScheduler { [ThreadStatic] private static bool _currentThreadIsProcessingItems; private readonly int _maxDegreeOfParallelism; private readonly LinkedList _tasks = new LinkedList(); // protected by lock(_tasks) private int _delegatesQueuedOrRunning = 0; // protected by lock(_tasks) /// Initializes an instance of the LimitedConcurrencyLevelTaskScheduler class with the specified degree of parallelism. /// The maximum degree of parallelism provided by this scheduler. public LimitedConcurrencyTaskScheduler (int maxDegreeOfParallelism) { if (maxDegreeOfParallelism < 1) throw new ArgumentOutOfRangeException(nameof(maxDegreeOfParallelism)); _maxDegreeOfParallelism = maxDegreeOfParallelism; } /// Gets the maximum concurrency level supported by this scheduler. public override sealed int MaximumConcurrencyLevel => _maxDegreeOfParallelism; protected override sealed IEnumerable GetScheduledTasks () { bool lockTaken = false; try { Monitor.TryEnter(_tasks, ref lockTaken); if (lockTaken) return _tasks.ToArray(); else throw new NotSupportedException(); } finally { if (lockTaken) Monitor.Exit(_tasks); } } protected override sealed void QueueTask (Task task) { // Add the task to the list of tasks to be processed. If there aren't enough // delegates currently queued or running to process tasks, schedule another. lock (_tasks) { _tasks.AddLast(task); if (_delegatesQueuedOrRunning < _maxDegreeOfParallelism) { ++_delegatesQueuedOrRunning; NotifyThreadPoolOfPendingWork(); } } } /// Informs the ThreadPool that there's work to be executed for this scheduler. private void NotifyThreadPoolOfPendingWork () { ThreadPool.UnsafeQueueUserWorkItem(_ => { // Note that the current thread is now processing work items. // This is necessary to enable inlining of tasks into this thread. _currentThreadIsProcessingItems = true; try { // Process all available items in the queue. while (true) { Task item; lock (_tasks) { // When there are no more items to be processed, // note that we're done processing, and get out. if (_tasks.Count == 0) { --_delegatesQueuedOrRunning; break; } // Get the next item from the queue item = _tasks.First.Value; _tasks.RemoveFirst(); } // Execute the task we pulled out of the queue TryExecuteTask(item); } } finally { // We're done processing items on the current thread _currentThreadIsProcessingItems = false; } }, null); } protected override sealed bool TryExecuteTaskInline (Task task, bool taskWasPreviouslyQueued) { // If this thread isn't already processing a task, we don't support inlining if (!_currentThreadIsProcessingItems) return false; // If the task was previously queued, remove it from the queue if (taskWasPreviouslyQueued) TryDequeue(task); // Try to run the task. return TryExecuteTask(task); } protected override sealed bool TryDequeue (Task task) { lock (_tasks) return _tasks.Remove(task); } } Ещё можно воспользоваться BlockingCollection<> с ConcurrentQueue<> внутри. По окончанию добавления надо будет не забыть вызвать CompleteAdding. Ещё есть TPL Dataflow. Там параллелизмы и прочие настраиваются без кастомных планировщиков. Можно ещё Rx добавить в качестве вишенки на торте. Вариантов море. Вообще, сейчас ещё @Vlad прибежит, расскажет про хитрые мозговыносящие способы решить эту задачу. :)

Ответ 2



Вам нужен MailboxProcessor - агент для обработки сообщений, который выполняет асинхронные операции. Реализация паттерна "много писателей, один читатель".

суббота, 13 июля 2019 г.

Promises. Не работает последовательность

Здравствуйте. Пишу код с использованием Обещаний для последовательного запуска функций. Только пару дней как разбираюсь с ними для улучшения кода (ибо раньше была ёлочка из коллбэков); Так вот. Есть функция анимации вывода текста на экран Animal():
function Animal(string) { var a = ''; //variable which will be entered character by character string var i= 0;//Letter counter var p = document.createElement('p'); $('body').scrollTop($('body').append(p).height()); return new Promise(function (resolve) { Anima(); function Anima() {//Function of Animation a+=string[i]; i++; $('p').last().text(a); var timer = setTimeout(Anima, 100); if(i==string.length){ clearTimeout(timer); resolve(); }; }; }); };
И функция AnimalPause() для изображения пауз(выводится на экран строка из точек и сразу же удаляется):
function AnimalPause(string) { var a = ''; //variable which will be entered character by character string var i= 0;//Letter counter var p = document.createElement('p'); $('body').scrollTop($('body').append(p).height()); return promise = new Promise(function (resolve) { Anima(); function Anima() {//Function of Animation a+=string[i]; i++; $('p').last().text(a); var timer = setTimeout(Anima, 100); if(i==string.length){ clearTimeout(timer); $('p').last().remove(); resolve(); }; }; }); };
Все сделал, вроде бы. И, если выводить последовательно строки через .then() - всё работает. И, если выводить последовательно строки с использованием цикла - тоже всё работает:
Animal('..........') .then(() => Animal('---Hello! CONSOLE v 1.0.1 is working!---')) .then(() => Animal('To see all commands u can use type -help')) .then(() => Animal('DONE'))
или так:
var chain = Promise.resolve(); pause.forEach(function(txt){ chain = chain.then(() => AnimalPause(txt))});
var pause = [ '...', '.....', '....', '.........' ];
Но если через .then() Последовательно выводить текст, потом задержку, потом - текст, рушится: (Задержка выводится после того, как выведется текст):
function Hello() { var chain = Animal('..........') .then(() => Animal('---Hello! CONSOLE v 1.0.1 is working!---')) .then(() => Animal('To see all commands u can use type -help')) .then(function() { return Pause(chain); // pause.forEach(function(txt){ // chain = chain.then(() => AnimalPause(txt))}) }) .then(() => Animal('DONE'))
}; Hello();
function Pause (chain) { return new Promise(function(resolve) { pause.forEach(function(txt){ chain = chain.then(() => AnimalPause(txt))}); resolve(); }) }
Почему? Ну..и как исправить? :)
По совету @Grundy в комментарии скидываю код для тестирования
var pause = [ '...', '.....', '....', '.........' ]; function Animal(string) { var a = ''; //variable which will be entered character by character string var i = 0; //Letter counter var p = document.createElement('p'); $('body').scrollTop($('body').append(p).height()); return new Promise(function(resolve) { Anima(); function Anima() { //Function of Animation a += string[i]; i++; $('p').last().text(a); var timer = setTimeout(Anima, 100); if (i == string.length) { clearTimeout(timer); resolve(); }; }; }); }; function AnimalPause(string) { var a = ''; //variable which will be entered character by character string var i = 0; //Letter counter var p = document.createElement('p'); $('body').scrollTop($('body').append(p).height()); return promise = new Promise(function(resolve) { Anima(); function Anima() { //Function of Animation a += string[i]; i++; $('p').last().text(a); var timer = setTimeout(Anima, 100); if (i == string.length) { clearTimeout(timer); $('p').last().remove(); resolve(); }; }; }); }; function Pause(chain) { return new Promise(function(resolve) { pause.forEach(function(txt) { chain = chain.then(() => AnimalPause(txt)) }); resolve(); }) }; function Hello() { var chain = Animal('..........') .then(() => Animal('---Hello! CONSOLE v 1.0.1 is working!---')) .then(() => Animal('To see all commands u can use type -help')) .then(function() { return Pause(chain); //Вы заметите, что строка пауз выводится после вывода на экран строки "DONE". Должно быть наоборот }) .then(() => Animal('DONE')) }; Hello();


Ответ

Если пройтись по коду, то можно отметить, что функция Animal отличается от AnimalPause только тем, что в последней в итоге удаляется добавленный элемент.
Исходя из этого можно передавать созданный p в resolve функции Animal и удалять его если надо. При этом AnimalPause выродится в следующее
function AnimalPause(string) { return Animal(string).then(p => p.remove()); };
Далее идет основная ошибка: функция Pause, которая добавляет продолжения для
var chain = Animal(...)
Но при этом не возвращает итоговый Promise, а просто переводит себя в состояние готово, именно поэтому выполнения вывода паузы откладывается до следующей цепочки.
Вместо этого нужно вернуть Promise собранный на основе массива pause с помощью функции reduce
function Pause(pause) { return pause.reduce((chain, txt) => chain.then(() => AnimalPause(txt)), Promise.resolve()); };
В этом случае возвращенный Promise будет встроен в существующую цепочку и вызван в нужном порядке.

Пример в сборе:
var pause = [ '...', '.....', '....', '.........' ]; function Animal(string) { var a = ''; //variable which will be entered character by character string var i = 0; //Letter counter var p = document.createElement('p'); $('body').scrollTop($('body').append(p).height()); return new Promise(function(resolve) { Anima(); function Anima() { //Function of Animation a += string[i]; i++; p.textContent = a; var timer = setTimeout(Anima, 100); if (i == string.length) { clearTimeout(timer); resolve(p); }; }; }); }; function AnimalPause(string) { return Animal(string).then(p => p.remove()); }; function Pause(pause) { return pause.reduce((chain, txt) => chain.then(() => AnimalPause(txt)), Promise.resolve()); }; function Hello() { var chain = Animal('..........') .then(() => Animal('---Hello! CONSOLE v 1.0.1 is working!---')) .then(() => Animal('To see all commands u can use type -help')) .then(() => Pause(pause)) .then(() => Animal('DONE')) }; Hello();

понедельник, 8 июля 2019 г.

Застревают очереди при отправке писем

Застревают очереди при отправки чере queue, если ставлю просто через send то отрабатывает стабильно
Mail::to($usermail, $username)->send(new \App\Mail\ConfirmEmail
Какие варианты решения этой проблемы существуют. Заранее благодарен. Почтовый сервер exim
Проблема в том, что это задерживает отработку скрипта так как очередь работает прозрачно и на отработку скрипта никак не влияет, а при прямой отправке происходит задержка и довольно существенная так как отправка идет не одному адресату, а десятку или даже сотне
Драйвер Redis


Ответ

При использовании QUEUE_DRIVER=database
Если создана таблица для сбора заданий завершенных с ошибками, там будет содержаться подробная информация. Если не создана таблица выполните
./artisan queue:failed-table ./artisan migrate
Снова запустите listener
./artisan queue:work --queue=<название_очереди> --sleep=2 --tries=1 --timeout 30 --daemon
Выполните отправку сообщения и посмотрите есть ли не выполненные задания
./artisan queue:failed
Если будет такой вывод
+----+------------+-------+-----------------------+--------------------- + | ID | Connection | Queue | Class | Failed At | +----+------------+-------+-----------------------+--------------------- + | 10 | database | email | App\Jobs\SendEmailJob | 2017-04-11 11:43:16 | | 9 | database | email | App\Jobs\SendEmailJob | 2017-04-11 11:40:16 |
выполните SQL запрос, в поле exception будет содержаться информация об ошибке.
SELECT * FROM failed_jobs;

четверг, 2 мая 2019 г.

Будет ли равномерным распределение заданий в очереди?

Нужно что-то вроде RR.
Будет ли равномерным распределение заданий в очереди Queue.Queue() по тредам? Смущает то что задания появляться будут медленней чем выполнение команд по ним.
Запускаю несколько тредов (с разным кодом под разные устройства) и тред который принимает задания по http. Задание попадает в очередь, а до этого воркеры висят заблокированном queue.get().
Я использовал в другом проекте multiprocessing.dummy.Pool - там все красиво и равномерно, но сейчас треды имеют разный код внутри и одной функцией не обойдешься.
Задача разгрузить исполнительные устройсва, а не процессор.


Ответ

Сперва почему-то хотелось ответить, что распределение неравномерное - то есть часть потоков будет простаивать, а тот, что был создан первее всех будет отдуваться за всех. Видимо такой поспешный вывод из-за наивного представления потоков в воображении. Но вот эксперимент (py3) - задания появляются в два раза медленнее, чем обрабатываются:
import queue import threading import time import random
counter = {} lock = threading.Lock()
def worker(): while True: try: sleep = q.get() if sleep == "DIE": q.task_done() return print(threading.current_thread().name) with lock: counter[threading.current_thread().name] = counter.get(threading.current_thread().name, 0) + 1 time.sleep(sleep) q.task_done() except queue.Empty: pass
q = queue.Queue() pool = [threading.Thread(target=worker, name="Worker " + str(i)) for i in range(4)] [t.start() for t in pool] for i in range(100): rand_sleep = random.random() / 8 q.put(rand_sleep) time.sleep(rand_sleep * 2) for i in range(4): q.put("DIE")
q.join()
print("Report: ", counter)
Вот что скрипт выводит:
output: ('Report: ', {'Worker 2': 25, 'Worker 3': 25, 'Worker 0': 25, 'Worker 1': 25})
Как видно из результатов - для моего кода распределение равномерное.
Если задания поступают без задержек - то распределение неравномерное и зависит от того, как долго обрабатывает задание каждый поток.
Дополнение: данный скрипт был протестирован в Windows и Ubuntu, использовался Python3.4 и Python2.7. Во всех четырех случаях результаты одинаковы.

среда, 5 декабря 2018 г.

Сервер очередей под задачу

Пытаюсь выбрать наиболее подходящий сервер очередей для такой задачи: Есть алгоритм, который выполняется для определенного интервала, нужно разбить этот интервал и вынести вычисления в воркеры (количество которых желательно задавать динамически в зависимости от размера интервала), далее дождаться выполнения всех воркеров и выполнить действия над вернувшимися значениями.
Я почитал, что для этого в принципе подходит gearman, но хотелось бы услышать еще и ваши мнения. Также читал про популярный rabbitMq, но не нагуглил в нем возможности задавать несколько воркеров
Язык программирования, в принципе, не важен, главное чтоб многопоточность поддерживал (многоядерную).


Ответ

RabbitMQ позволит тебе подключить к одной очереди несколько consumer-ов (см. http://www.rabbitmq.com/tutorials/tutorial-two-python.html Fair Dispatch). По поводу динамических консьюмеров - тут твое приложение само должно решать когда их запускать и как их выключать. Можешь запустить кол-во консьюмеров равное кол-ву ядер на машине и потом собирать данные.
А твоя задача напоминает реализацию MapReduce. Может тебе имеет смысл взглянуть на Hadoop?
Если все же хочешь написать сам, то по поводу языка можно выбрать Erlang/OTP (на нем написан RabbitMQ). Работа с потоками там построена очень хорошо и объединить несколько машин в кластер, с шарингом потоков, не составит труда. Твое приложение не будет думать о том, где и что выполняется, а просто ждать завершения всех воркеров и вернет результат. Но сломать мозг придется, прежде чем написать функциональный код :)

вторник, 13 ноября 2018 г.

BlockingCollection - как не блокировать поток

Применение BlockingCollection, используя подход, когда элементы вытаскиваются из очереди(например ConcurrentQueue), используя метод Take в цикле - всегда блокирует поток. Очевидно, что процессорное время не занимается, однако поток все же занят и не может использоваться для выполнения других задач. Какая есть альтернатива, когда нужно последовательно вычитывать элементы из очереди и при этом не блокировать поток? Конечно, можно сделать велосипед, накрутить событий или чего-нибудь еще, но хотелось бы понять, нет ли каких-либо стандартных способов это сделать, кроме как использовать BlockingCollection и метод Take
P.S. Есть метод TryTake, но я не могу найти решение с его использованием, эквивалентное использованию Take и при этом неблокирующее поток.


Ответ

Вам на самом деле нужен класс BufferBlock из библиотеки Dataflow (nuget-пакет Microsoft.Tpl.Dataflow).
Этот класс заменяет собой BlockingCollection, и позволяет асинхронный доступ:
await queue.ReceiveAsync()
Таким образом, поток не будет заблокирован. Но у вас получится async-интерфейс.
Больше примеров с работающим кодом есть в этом ответе

Ещё одним вариантом является async-обёртка над IProducerCosumerCollection из AsyncEx Стивена Клири: https://github.com/StephenCleary/AsyncEx/wiki/AsyncCollection

четверг, 4 октября 2018 г.

Ревью класса для очереди команд

Мне потребовалась очередь команд, где множество потоков может добавлять команды на выполнение и один поток по очереди их выполняет. Тк я не знаю какие-либо стандартные реализации подобного, то пришлось делать самому. Можете оценить мой класс для этого, тк у меня сомнения на счет его качества, хотя он работает. И если кто-либо даст ссылку на реализации подобного, то тоже будет приятно.
public class ConcurrentCommandsQueue : IDisposable { private readonly object _lockForSyncExecuting; private readonly Queue _commandQueue;
private readonly Action _executeCommand; private readonly Action _onStartCommandsExecution; private readonly Action _onEndCommandsExecution;
private volatile bool _isDisposed; private volatile bool _isExecuting; private volatile Task _taskForAction;
public bool IsExecuting { get { return _isExecuting; } }
public ConcurrentCommandsQueue(Action executeCommand) : this(executeCommand, null, null) { } public ConcurrentCommandsQueue(Action executeCommand, Action onStartCommandsExecution, Action onEndCommandsExecution) { _lockForSyncExecuting = new object(); _commandQueue = new Queue(10);
_isDisposed = false; _isExecuting = false; _executeCommand = executeCommand; _onStartCommandsExecution = onStartCommandsExecution; _onEndCommandsExecution = onEndCommandsExecution; _taskForAction = new TaskCompletionSource().Task; }
private void ExecuteCommands() { if (_onStartCommandsExecution != null) _onStartCommandsExecution(); T command = default(T); while (true) { lock (_lockForSyncExecuting) { if (_commandQueue.Count == 0) { _isExecuting = false; break; } command = _commandQueue.Dequeue(); } _executeCommand(command); } if (_onEndCommandsExecution != null) _onEndCommandsExecution(command); }
public void AddCommand(T command) { if (_isDisposed) throw new ObjectDisposedException(GetType().ToString());
lock (_lockForSyncExecuting) { if (_isDisposed) return; _commandQueue.Enqueue(command); if (!_isExecuting) { _isExecuting = true; _taskForAction = Task.Run(new Action(ExecuteCommands)); } } } public void ClearCommandQueue() { if (_isDisposed) throw new ObjectDisposedException(GetType().ToString());
lock (_lockForSyncExecuting) _commandQueue.Clear(); } public void Dispose() { if (_isDisposed) return;
_isDisposed = true; lock (_lockForSyncExecuting) _commandQueue.Clear(); if (_taskForAction != null) _taskForAction.Wait(); } }


Ответ

У вас получился код в стиле .NET 1.0, который слегка эволюционировал с приходом лямбд (Delegate механически заменён на Action) и задач (new Thread().Start() механически заменён на Task.Run()), но по сути совершенно не изменился: вы TPL с async/await толком не используете, а существование concurrent коллекций вообще упустили из виду.
В .NET 4.5 это всё ненужно в принципе. Если вам нужно выполнение задач строго в одном потоке, то можно просто кидать задачи в однопоточный планировщик задач (task scheduler). Собственно, на этой первой строчке весь код и заканчивается. Если нужно выполнение кода перед всеми операциями и после всех операций, то код перед и после этой строчки и пишется. Отмечу, что при этой однострочной реализации вы ещё имеете бонусом возможность использования токена отмены (cancellation token) и прочих радостей жизни.

Рассмотрим пример. Допустим, надо выполнить две очереди задач: в каждой очереди задачи выполняются последовательно, но очереди выполняются параллельно. Создаём два планировщика задач с ограничением параллелизма, затем создаём задачи с указанием нужного нам планировщика.
class Program { static readonly Random _rnd = new Random(); static readonly LimitedConcurrencyTaskScheduler _schedulerFoo = new LimitedConcurrencyTaskScheduler(1); static readonly LimitedConcurrencyTaskScheduler _schedulerBar = new LimitedConcurrencyTaskScheduler(1);
static void Main () => new Program().Run().Wait();
async Task Run () { Task queueFoo = RunQueue("Foo", _schedulerFoo, Enumerable.Range(0, 3).Select(i => (Action)(() => Foo("Foo")))); Task queueBar = RunQueue("Bar", _schedulerBar, Enumerable.Range(0, 3).Select(i => (Action)(() => Foo("Bar")))); await Task.WhenAll(queueFoo, queueBar); Console.WriteLine("Done!"); Console.ReadKey(); }
async Task RunQueue (string name, TaskScheduler scheduler, IEnumerable commands) { Console.WriteLine($"{name}: Start"); await Task.WhenAll(commands.Select(c => RunTask(c, scheduler))); Console.WriteLine($"{name}: Finish"); }
async Task RunTask (Action task, TaskScheduler scheduler) { await Task.Factory.StartNew(task, CancellationToken.None, TaskCreationOptions.None, scheduler); }
void Foo (string name) { int timeout = _rnd.Next(200); Console.WriteLine($"{name}: Start {timeout}"); Thread.Sleep(timeout); Console.WriteLine($"{name}: Finish {timeout}"); } }
Пример вывода:
Foo: Start Bar: Start Foo: Start 165 Bar: Start 50 Bar: Finish 50 Bar: Start 39 Bar: Finish 39 Bar: Start 115 Foo: Finish 165 Foo: Start 116 Bar: Finish 115 Bar: Finish Foo: Finish 116 Foo: Start 0 Foo: Finish 0 Foo: Finish Done!
Хотя в этом примере все задачи и кидаются одновременно, задачи можно добавлять в любой момент. Правда тогда смысл "начала" и "конца" несколько теряется. Если вы объясните, что вы собираетесь с ними делать, тогда можно будет добавить.
Здесь я воспользовался планировщиком LimitedConcurrencyTaskScheduler из примеров на MSDN: Samples for Parallel Programming with the .NET Framework (статья с описанием: ParallelExtensionsExtras Tour - #7 - Additional TaskSchedulers).
///

/// Provides a task scheduler that ensures a maximum concurrency level while running on top of the ThreadPool. /// Source: http://code.msdn.microsoft.com/ParExtSamples /// Documentation: http://blogs.msdn.com/b/pfxteam/archive/2010/04/09/9990424.aspx /// License: MS-LPL /// public class LimitedConcurrencyTaskScheduler : TaskScheduler { [ThreadStatic] private static bool _currentThreadIsProcessingItems;
private readonly int _maxDegreeOfParallelism; private readonly LinkedList _tasks = new LinkedList(); // protected by lock(_tasks) private int _delegatesQueuedOrRunning = 0; // protected by lock(_tasks)
/// Initializes an instance of the LimitedConcurrencyLevelTaskScheduler class with the specified degree of parallelism. /// The maximum degree of parallelism provided by this scheduler. public LimitedConcurrencyTaskScheduler (int maxDegreeOfParallelism) { if (maxDegreeOfParallelism < 1) throw new ArgumentOutOfRangeException(nameof(maxDegreeOfParallelism)); _maxDegreeOfParallelism = maxDegreeOfParallelism; }
/// Gets the maximum concurrency level supported by this scheduler. public override sealed int MaximumConcurrencyLevel => _maxDegreeOfParallelism;
protected override sealed IEnumerable GetScheduledTasks () { bool lockTaken = false; try { Monitor.TryEnter(_tasks, ref lockTaken); if (lockTaken) return _tasks.ToArray(); else throw new NotSupportedException(); } finally { if (lockTaken) Monitor.Exit(_tasks); } }
protected override sealed void QueueTask (Task task) { // Add the task to the list of tasks to be processed. If there aren't enough // delegates currently queued or running to process tasks, schedule another. lock (_tasks) { _tasks.AddLast(task); if (_delegatesQueuedOrRunning < _maxDegreeOfParallelism) { ++_delegatesQueuedOrRunning; NotifyThreadPoolOfPendingWork(); } } }
/// Informs the ThreadPool that there's work to be executed for this scheduler. private void NotifyThreadPoolOfPendingWork () { ThreadPool.UnsafeQueueUserWorkItem(_ => { // Note that the current thread is now processing work items. // This is necessary to enable inlining of tasks into this thread. _currentThreadIsProcessingItems = true; try { // Process all available items in the queue. while (true) { Task item; lock (_tasks) { // When there are no more items to be processed, // note that we're done processing, and get out. if (_tasks.Count == 0) { --_delegatesQueuedOrRunning; break; } // Get the next item from the queue item = _tasks.First.Value; _tasks.RemoveFirst(); } // Execute the task we pulled out of the queue TryExecuteTask(item); } } finally { // We're done processing items on the current thread _currentThreadIsProcessingItems = false; } }, null); }
protected override sealed bool TryExecuteTaskInline (Task task, bool taskWasPreviouslyQueued) { // If this thread isn't already processing a task, we don't support inlining if (!_currentThreadIsProcessingItems) return false; // If the task was previously queued, remove it from the queue if (taskWasPreviouslyQueued) TryDequeue(task); // Try to run the task. return TryExecuteTask(task); }
protected override sealed bool TryDequeue (Task task) { lock (_tasks) return _tasks.Remove(task); } }

Ещё можно воспользоваться BlockingCollection<> с ConcurrentQueue<> внутри. По окончанию добавления надо будет не забыть вызвать CompleteAdding
Ещё есть TPL Dataflow. Там параллелизмы и прочие настраиваются без кастомных планировщиков.
Можно ещё Rx добавить в качестве вишенки на торте.
Вариантов море.
Вообще, сейчас ещё @Vlad прибежит, расскажет про хитрые мозговыносящие способы решить эту задачу. :)

Ревью класса для очереди команд

Мне потребовалась очередь команд, где множество потоков может добавлять команды на выполнение и один поток по очереди их выполняет. Тк я не знаю какие-либо стандартные реализации подобного, то пришлось делать самому. Можете оценить мой класс для этого, тк у меня сомнения на счет его качества, хотя он работает. И если кто-либо даст ссылку на реализации подобного, то тоже будет приятно.
public class ConcurrentCommandsQueue : IDisposable { private readonly object _lockForSyncExecuting; private readonly Queue _commandQueue;
private readonly Action _executeCommand; private readonly Action _onStartCommandsExecution; private readonly Action _onEndCommandsExecution;
private volatile bool _isDisposed; private volatile bool _isExecuting; private volatile Task _taskForAction;
public bool IsExecuting { get { return _isExecuting; } }
public ConcurrentCommandsQueue(Action executeCommand) : this(executeCommand, null, null) { } public ConcurrentCommandsQueue(Action executeCommand, Action onStartCommandsExecution, Action onEndCommandsExecution) { _lockForSyncExecuting = new object(); _commandQueue = new Queue(10);
_isDisposed = false; _isExecuting = false; _executeCommand = executeCommand; _onStartCommandsExecution = onStartCommandsExecution; _onEndCommandsExecution = onEndCommandsExecution; _taskForAction = new TaskCompletionSource().Task; }
private void ExecuteCommands() { if (_onStartCommandsExecution != null) _onStartCommandsExecution(); T command = default(T); while (true) { lock (_lockForSyncExecuting) { if (_commandQueue.Count == 0) { _isExecuting = false; break; } command = _commandQueue.Dequeue(); } _executeCommand(command); } if (_onEndCommandsExecution != null) _onEndCommandsExecution(command); }
public void AddCommand(T command) { if (_isDisposed) throw new ObjectDisposedException(GetType().ToString());
lock (_lockForSyncExecuting) { if (_isDisposed) return; _commandQueue.Enqueue(command); if (!_isExecuting) { _isExecuting = true; _taskForAction = Task.Run(new Action(ExecuteCommands)); } } } public void ClearCommandQueue() { if (_isDisposed) throw new ObjectDisposedException(GetType().ToString());
lock (_lockForSyncExecuting) _commandQueue.Clear(); } public void Dispose() { if (_isDisposed) return;
_isDisposed = true; lock (_lockForSyncExecuting) _commandQueue.Clear(); if (_taskForAction != null) _taskForAction.Wait(); } }


Ответ

У вас получился код в стиле .NET 1.0, который слегка эволюционировал с приходом лямбд (Delegate механически заменён на Action) и задач (new Thread().Start() механически заменён на Task.Run()), но по сути совершенно не изменился: вы TPL с async/await толком не используете, а существование concurrent коллекций вообще упустили из виду.
В .NET 4.5 это всё ненужно в принципе. Если вам нужно выполнение задач строго в одном потоке, то можно просто кидать задачи в однопоточный планировщик задач (task scheduler). Собственно, на этой первой строчке весь код и заканчивается. Если нужно выполнение кода перед всеми операциями и после всех операций, то код перед и после этой строчки и пишется. Отмечу, что при этой однострочной реализации вы ещё имеете бонусом возможность использования токена отмены (cancellation token) и прочих радостей жизни.

Рассмотрим пример. Допустим, надо выполнить две очереди задач: в каждой очереди задачи выполняются последовательно, но очереди выполняются параллельно. Создаём два планировщика задач с ограничением параллелизма, затем создаём задачи с указанием нужного нам планировщика.
class Program { static readonly Random _rnd = new Random(); static readonly LimitedConcurrencyTaskScheduler _schedulerFoo = new LimitedConcurrencyTaskScheduler(1); static readonly LimitedConcurrencyTaskScheduler _schedulerBar = new LimitedConcurrencyTaskScheduler(1);
static void Main () => new Program().Run().Wait();
async Task Run () { Task queueFoo = RunQueue("Foo", _schedulerFoo, Enumerable.Range(0, 3).Select(i => (Action)(() => Foo("Foo")))); Task queueBar = RunQueue("Bar", _schedulerBar, Enumerable.Range(0, 3).Select(i => (Action)(() => Foo("Bar")))); await Task.WhenAll(queueFoo, queueBar); Console.WriteLine("Done!"); Console.ReadKey(); }
async Task RunQueue (string name, TaskScheduler scheduler, IEnumerable commands) { Console.WriteLine($"{name}: Start"); await Task.WhenAll(commands.Select(c => RunTask(c, scheduler))); Console.WriteLine($"{name}: Finish"); }
async Task RunTask (Action task, TaskScheduler scheduler) { await Task.Factory.StartNew(task, CancellationToken.None, TaskCreationOptions.None, scheduler); }
void Foo (string name) { int timeout = _rnd.Next(200); Console.WriteLine($"{name}: Start {timeout}"); Thread.Sleep(timeout); Console.WriteLine($"{name}: Finish {timeout}"); } }
Пример вывода:
Foo: Start Bar: Start Foo: Start 165 Bar: Start 50 Bar: Finish 50 Bar: Start 39 Bar: Finish 39 Bar: Start 115 Foo: Finish 165 Foo: Start 116 Bar: Finish 115 Bar: Finish Foo: Finish 116 Foo: Start 0 Foo: Finish 0 Foo: Finish Done!
Хотя в этом примере все задачи и кидаются одновременно, задачи можно добавлять в любой момент. Правда тогда смысл "начала" и "конца" несколько теряется. Если вы объясните, что вы собираетесь с ними делать, тогда можно будет добавить.
Здесь я воспользовался планировщиком LimitedConcurrencyTaskScheduler из примеров на MSDN: Samples for Parallel Programming with the .NET Framework (статья с описанием: ParallelExtensionsExtras Tour - #7 - Additional TaskSchedulers).
///

/// Provides a task scheduler that ensures a maximum concurrency level while running on top of the ThreadPool. /// Source: http://code.msdn.microsoft.com/ParExtSamples /// Documentation: http://blogs.msdn.com/b/pfxteam/archive/2010/04/09/9990424.aspx /// License: MS-LPL /// public class LimitedConcurrencyTaskScheduler : TaskScheduler { [ThreadStatic] private static bool _currentThreadIsProcessingItems;
private readonly int _maxDegreeOfParallelism; private readonly LinkedList _tasks = new LinkedList(); // protected by lock(_tasks) private int _delegatesQueuedOrRunning = 0; // protected by lock(_tasks)
/// Initializes an instance of the LimitedConcurrencyLevelTaskScheduler class with the specified degree of parallelism. /// The maximum degree of parallelism provided by this scheduler. public LimitedConcurrencyTaskScheduler (int maxDegreeOfParallelism) { if (maxDegreeOfParallelism < 1) throw new ArgumentOutOfRangeException(nameof(maxDegreeOfParallelism)); _maxDegreeOfParallelism = maxDegreeOfParallelism; }
/// Gets the maximum concurrency level supported by this scheduler. public override sealed int MaximumConcurrencyLevel => _maxDegreeOfParallelism;
protected override sealed IEnumerable GetScheduledTasks () { bool lockTaken = false; try { Monitor.TryEnter(_tasks, ref lockTaken); if (lockTaken) return _tasks.ToArray(); else throw new NotSupportedException(); } finally { if (lockTaken) Monitor.Exit(_tasks); } }
protected override sealed void QueueTask (Task task) { // Add the task to the list of tasks to be processed. If there aren't enough // delegates currently queued or running to process tasks, schedule another. lock (_tasks) { _tasks.AddLast(task); if (_delegatesQueuedOrRunning < _maxDegreeOfParallelism) { ++_delegatesQueuedOrRunning; NotifyThreadPoolOfPendingWork(); } } }
/// Informs the ThreadPool that there's work to be executed for this scheduler. private void NotifyThreadPoolOfPendingWork () { ThreadPool.UnsafeQueueUserWorkItem(_ => { // Note that the current thread is now processing work items. // This is necessary to enable inlining of tasks into this thread. _currentThreadIsProcessingItems = true; try { // Process all available items in the queue. while (true) { Task item; lock (_tasks) { // When there are no more items to be processed, // note that we're done processing, and get out. if (_tasks.Count == 0) { --_delegatesQueuedOrRunning; break; } // Get the next item from the queue item = _tasks.First.Value; _tasks.RemoveFirst(); } // Execute the task we pulled out of the queue TryExecuteTask(item); } } finally { // We're done processing items on the current thread _currentThreadIsProcessingItems = false; } }, null); }
protected override sealed bool TryExecuteTaskInline (Task task, bool taskWasPreviouslyQueued) { // If this thread isn't already processing a task, we don't support inlining if (!_currentThreadIsProcessingItems) return false; // If the task was previously queued, remove it from the queue if (taskWasPreviouslyQueued) TryDequeue(task); // Try to run the task. return TryExecuteTask(task); }
protected override sealed bool TryDequeue (Task task) { lock (_tasks) return _tasks.Remove(task); } }

Ещё можно воспользоваться BlockingCollection<> с ConcurrentQueue<> внутри. По окончанию добавления надо будет не забыть вызвать CompleteAdding
Ещё есть TPL Dataflow. Там параллелизмы и прочие настраиваются без кастомных планировщиков.
Можно ещё Rx добавить в качестве вишенки на торте.
Вариантов море.
Вообще, сейчас ещё @Vlad прибежит, расскажет про хитрые мозговыносящие способы решить эту задачу. :)