Страницы

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

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

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

Организация доступа к ресурсу из разных потоков

#c_sharp #многопоточность #backgroundworker


Имеется такой код.

// Коллекция буферов на запись
List> writeReadyList = new List>();

private void backgroundWorker1_DoWork(object sender, DoWorkEventArgs e)     
{
    // Локальный буфер данного потока
    List localBufer = new List();
    Random random = new Random();

    while (!backgroundWorker1.CancellationPending)
    {
        // Забиваем локальный буфер значениями
        localBufer.Add(random.Next(0, 100));

        if (localBufer.Count > 10)            
        {
            // Передаем буфер в очередь на запись и обнуляем его
            writeReadyList.Add(localBufer);
            localBufer = new List();
        }

        Thread.Sleep(50);
    }
}

private void backgroundWorker2_DoWork(object sender, DoWorkEventArgs e)
{
    while (!backgroundWorker2.CancellationPending) 
    {
        // Коллекция буферов, которые не удалось записать
        List> failList = new List>(); 

        foreach (List currentBufer in writeReadyList)
        {
            try
            {
                // Пишем значения из перебираемого буфера в базу данных
            }
            catch { failList.Add(currentBufer); }
        }

        writeReadyList = new List>(failList);

        Thread.Sleep(450);
    }
}


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

Насколько понимаю, при таком подходе возможны несколько вариантов некорректной работы.


в момент работы foreach во втором потоке она будет изменена из первого потока - получаем
исключение?
второй поток выполнил перебор очереди, но еще не выполнил перезаписи в соответствии
с неудавшимися попытками записи. В этот момент первый добавляет запись, но, так как
второй ничего об этом не знает, то запись будет потеряна.


Как правильно организовать работу с очередью на запись между этими потоками? Причем
организовать надо так, чтобы если второй поток в данный момент обрабатывает очередь,
то первый поток не ждал окончания работы с ресурсом, а продолжал складывать значения
в свой локальный буфер. В какую сторону копать?

______UPD:

По совету VladD использовал мьютексы.

    Mutex writeReadyAccess = new Mutex();

    private void backgroundWorker1_DoWork(object sender, DoWorkEventArgs e)
    {
        if (localBufer.Count > 10)
            // Таймаут в 1 секунду тк нельзя тормозить первый поток на время обработки
очереди
            if (writeReadyAccess.WaitOne(1))
            {
                // Если второй поток не обрабатывает очередь, то добавляем в нее
локальный буфер и обнуляем его
                writeReadyAccess.ReleaseMutex();
            }
        // Если очередь была в обработке на момент запроса, то продолжаем писать
в локальный буфер        
    }

    private void backgroundWorker2_DoWork(object sender, DoWorkEventArgs e)
    {
        // Этот поток можно тормозить по времени, поэтому ждем без таймаута
        if (writeReadyAccess.WaitOne())
        {
            // Обработка очереди, в том числе ее перезапись
            writeReadyAccess.ReleaseMutex();
        }
    }

    


Ответы

Ответ 1



using System.Collections.Concurrent; using System.Threading; class Dto { public int State; public int Index; } class Test { BlockingCollection list = new BlockingCollection(); public void Run() { Task.Run(() => backgroundWorker2()); Task.Run(() => backgroundWorker1()); } void backgroundWorker1() { for(int i=0; i < 3; i++) list.Add(new Dto() { Index = i }); } void backgroundWorker2() { foreach (var v in list.GetConsumingEnumerable()) { // сохраняем данные Console.WriteLine("index=" + v.Index + " state=" + v.State); // если сохранить не удалось, то возвращаем в list if (v.Index == 0 && v.State++ < 2) list.Add(v); if(list.Count == 0) break; } Console.WriteLine("completed"); } } var t = new Test(); t.Run(); Результат index=0 state=0 index=1 state=0 index=2 state=0 index=0 state=1 вторая попытка записи index=0 state=2 еще одна completed Если надо передавать данные порциями, то определите коллекцию, например, так BlockingCollection UPDATE Другая версия, в которой backgroundWorker1 работает медленно class Dto { public int Index; } class Test { BlockingCollection list = new BlockingCollection(); public void Run() { Task.Run(() => backgroundWorker2()); Task.Run(() => backgroundWorker1()); } void backgroundWorker1() { for (int i = 0; i < 3; i++) { list.Add(new Dto() { Index = i }); // Thread.Sleep(1000); } list.CompleteAdding(); } void backgroundWorker2() { var retry = new Queue(); foreach (var v in list.GetConsumingEnumerable()) { // сохраняем данные Console.WriteLine("index=" + v.Index); // если сохранить не удалось, то ставим в очередь if (v.Index == 0) retry.Enqueue(v); } // повторяем попытку записи while(retry.Count > 0) Console.WriteLine("retry=" + retry.Dequeue().Index); Console.WriteLine("completed"); } } var t = new Test(); t.Run(); Результат index=0 index=1 index=2 retry=0 повторная попытка записи completed UPDATE Для конвейеризации обработки данных в разных потоках предназначена библиотека потоков данных - TPL Dataflow. Описание с примерами есть в MSDN и есть nuget-пакет Microsoft TPL Dataflow.

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

Работа с Backgroundworker и Dispatcher

#c_sharp #многопоточность #mvvm #backgroundworker


Пишу программу на C# с mvvm. У меня есть два эквивалентных куска кода, в которых,по-моему
мнению, должна происходить абсолютно одинаковая работа. Суть в чем- постепенная подгрузка(добавление)
элементов в коллекцию с помощью BackgroundWorker. Коллекция имеет биндинг с listview
 и соответственно  постепенное(пообъектное) добавление в коллекцию отображается в этом
listview.
Код:

public class ListContactViewModel : ViewModelBase
{
 public ObservableCollection DialogListFirstPage { get;set;}
 Dispatcher _dispatcher;
  public ListContactViewModel(VkApi vk)
    {  
        _dispatcher = Application.Current.Dispatcher;

       DialogListFirstPage = new ObservableCollection();
        var bw1 = new BackgroundWorker();
        bw1.DoWork += (o, e) =>
        {
            FirstDialogPageMethod(vk); //Создание коллекции диалогов

        };
        bw1.RunWorkerAsync();

 //внимания в этом методе достойна только строчка добавления в коллекцию
  private void FirstDialogPageMethod(VkApi vk)
    {
        int totalCount, unreadCount;
        var GetFirstDialigPage = vk.Messages.GetDialogs(20, 0, out totalCount, out
unreadCount);


        foreach (var i in GetFirstDialigPage)
        {

            var Names = vk.Users.Get(i.UserId.ToString(), ProfileFields.FirstName);

           _dispatcher.Invoke(()=> DialogListFirstPage.Add(new MessageChild() { AuthorFirstName
= Names.FirstName, AuthorLastName = Names.LastName, Body = i.Body, UserId = i.UserId,
Date = i.Date, ChatActiveIds = i.ChatActiveIds, Title = i.Title, UsersCount = i.UsersCount,
ChatId = i.ChatId }));

        }        
    }


на всякий случай приведу строчку биндинга из xaml:

 


И все закономерно: окно открывается пустым и я наблюдаю постоянное добавление элементов
в список.

Ситуация 2: Из предыдущего окна я перехожу в следующее , в котором аналогичная ситуация

  public class CurrentDialogViewModel:ViewModelBase
  {
   public ObservableCollection ReadyCollection { get; set; }
   Dispatcher disp;
     public CurrentDialogViewModel(VkApi vk,MessageChild parametr)
   {
       disp = Application.Current.Dispatcher;
        ReadyCollection = new ObservableCollection();
         var bw2 = new BackgroundWorker();
       bw2.DoWork += (p, m) =>
       {

           MoreMessages();

       };
       bw2.RunWorkerAsync();
 private void MoreMessages()
   {   
       foreach (var i in builder.ConcreateDialogCreater(ids))
       {
           disp.Invoke(() => ReadyCollection.Insert(0, i));

       }
       Datas.offset += 200;


   }


где builder.ConcreateDialogCreater(ids) возвращает     

 ObservableCollection 


Xaml:

   


Так вот в этом случае при открытии окна, оно у меня некоторое время остается пустым,
после чего список мгновенно отображает все объекты в ReadyCollection. 
От Insert это не зависит, с Add тоже самое. Также это не зависит от builder.ConcreateDialogCreater(ids),
потому что пробовал делать просто инициализацию объекта при добавлении в цикле

 ReadyCollection.Add(new MessageChild());


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

UPD: 
Продебажил еще раз - все таки я был не прав и задержка связана с выполнением  метода
builder.ConcreateDialogCreater(ids). И пока он не выполнится весь- foreach не начнется.
Ведь в первом случае я коллекцию заполняю непосредственно в том классе и задержка
обоснована работой библиотечных методов перед добавлением.
Во-втором же случае нужно ждать,пока метод выполнится  полностью.
    


Ответы

Ответ 1



В DoWork вместо Dispatcher используйте ReportProgress. (для его работы надо включить WorkerReportsProgress). partial class MainWindow : Window { public MainWindow() { this.DataContext = _List = new ObservableCollection(); } ObservableCollection _List; private void Button_Click(object sender, RoutedEventArgs e) { var w = new BackgroundWorker() { WorkerReportsProgress = true }; w.DoWork += (s, we) => { for (var i = 0; i < 50; i++) { Thread.Sleep(100); // тут что-то делаем w.ReportProgress(0, i); } }; w.ProgressChanged += (s, we) => // выполняется в основном потоке _List.Add(new Message() { Text = "t" + we.UserState }); w.RunWorkerAsync(); } class Message { public string Text { get; set; } } }

пятница, 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. Если хоститесь на амазоне, то с их лямбдами вроде как подобное можно реализовать.

четверг, 14 февраля 2019 г.

Работа с Backgroundworker и Dispatcher

Пишу программу на C# с mvvm. У меня есть два эквивалентных куска кода, в которых,по-моему мнению, должна происходить абсолютно одинаковая работа. Суть в чем- постепенная подгрузка(добавление) элементов в коллекцию с помощью BackgroundWorker. Коллекция имеет биндинг с listview и соответственно постепенное(пообъектное) добавление в коллекцию отображается в этом listview. Код:
public class ListContactViewModel : ViewModelBase { public ObservableCollection DialogListFirstPage { get;set;} Dispatcher _dispatcher; public ListContactViewModel(VkApi vk) { _dispatcher = Application.Current.Dispatcher;
DialogListFirstPage = new ObservableCollection(); var bw1 = new BackgroundWorker(); bw1.DoWork += (o, e) => { FirstDialogPageMethod(vk); //Создание коллекции диалогов
}; bw1.RunWorkerAsync();
//внимания в этом методе достойна только строчка добавления в коллекцию private void FirstDialogPageMethod(VkApi vk) { int totalCount, unreadCount; var GetFirstDialigPage = vk.Messages.GetDialogs(20, 0, out totalCount, out unreadCount);
foreach (var i in GetFirstDialigPage) {
var Names = vk.Users.Get(i.UserId.ToString(), ProfileFields.FirstName);
_dispatcher.Invoke(()=> DialogListFirstPage.Add(new MessageChild() { AuthorFirstName = Names.FirstName, AuthorLastName = Names.LastName, Body = i.Body, UserId = i.UserId, Date = i.Date, ChatActiveIds = i.ChatActiveIds, Title = i.Title, UsersCount = i.UsersCount, ChatId = i.ChatId }));
} }
на всякий случай приведу строчку биндинга из xaml:

И все закономерно: окно открывается пустым и я наблюдаю постоянное добавление элементов в список.
Ситуация 2: Из предыдущего окна я перехожу в следующее , в котором аналогичная ситуация
public class CurrentDialogViewModel:ViewModelBase { public ObservableCollection ReadyCollection { get; set; } Dispatcher disp; public CurrentDialogViewModel(VkApi vk,MessageChild parametr) { disp = Application.Current.Dispatcher; ReadyCollection = new ObservableCollection(); var bw2 = new BackgroundWorker(); bw2.DoWork += (p, m) => {
MoreMessages();
}; bw2.RunWorkerAsync(); private void MoreMessages() { foreach (var i in builder.ConcreateDialogCreater(ids)) { disp.Invoke(() => ReadyCollection.Insert(0, i));
} Datas.offset += 200;
}
где builder.ConcreateDialogCreater(ids) возвращает
ObservableCollection
Xaml:

Так вот в этом случае при открытии окна, оно у меня некоторое время остается пустым, после чего список мгновенно отображает все объекты в ReadyCollection. От Insert это не зависит, с Add тоже самое. Также это не зависит от builder.ConcreateDialogCreater(ids), потому что пробовал делать просто инициализацию объекта при добавлении в цикле
ReadyCollection.Add(new MessageChild());
Аналогичная история-объекты вываливаются всем скопом по окончании добавления последнего. А я хочу добиться постепенной подгрузки, как в предыдущем окне. Почему так происходит и что нужно исправить?
UPD: Продебажил еще раз - все таки я был не прав и задержка связана с выполнением метода builder.ConcreateDialogCreater(ids). И пока он не выполнится весь- foreach не начнется. Ведь в первом случае я коллекцию заполняю непосредственно в том классе и задержка обоснована работой библиотечных методов перед добавлением. Во-втором же случае нужно ждать,пока метод выполнится полностью.


Ответ

В DoWork вместо Dispatcher используйте ReportProgress. (для его работы надо включить WorkerReportsProgress).
partial class MainWindow : Window { public MainWindow() { this.DataContext = _List = new ObservableCollection(); }
ObservableCollection _List;
private void Button_Click(object sender, RoutedEventArgs e) { var w = new BackgroundWorker() { WorkerReportsProgress = true }; w.DoWork += (s, we) => { for (var i = 0; i < 50; i++) { Thread.Sleep(100); // тут что-то делаем w.ReportProgress(0, i); } }; w.ProgressChanged += (s, we) => // выполняется в основном потоке _List.Add(new Message() { Text = "t" + we.UserState }); w.RunWorkerAsync(); }
class Message { public string Text { get; set; } } }

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

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

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


Ответ

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