Страницы

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

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

суббота, 11 апреля 2020 г.

Как gunicorn запустить для обслуживания нескольких потоков?

#python #многопоточность #gunicorn

                    
Есть тестовое приложение:

import time
from flask import Flask

app = Flask(import_name=__name__)


@app.route('/')
def test():
    time.sleep(5)
    return 'ok'


Запускаю gunicorn так:

(VENV) p2mbot@mbp:~/projects/test/gt$ gunicorn --workers 10 --threads 10 gt:app
[2015-07-18 11:23:10 +0300] [13196] [INFO] Starting gunicorn 19.3.0
[2015-07-18 11:23:10 +0300] [13196] [INFO] Listening at: http://127.0.0.1:8000 (13196)
[2015-07-18 11:23:10 +0300] [13196] [INFO] Using worker: threads
[2015-07-18 11:23:10 +0300] [13199] [INFO] Booting worker with pid: 13199
[2015-07-18 11:23:10 +0300] [13200] [INFO] Booting worker with pid: 13200
[2015-07-18 11:23:10 +0300] [13201] [INFO] Booting worker with pid: 13201
[2015-07-18 11:23:10 +0300] [13202] [INFO] Booting worker with pid: 13202
[2015-07-18 11:23:10 +0300] [13203] [INFO] Booting worker with pid: 13203
[2015-07-18 11:23:10 +0300] [13204] [INFO] Booting worker with pid: 13204
[2015-07-18 11:23:10 +0300] [13205] [INFO] Booting worker with pid: 13205
[2015-07-18 11:23:10 +0300] [13206] [INFO] Booting worker with pid: 13206
[2015-07-18 11:23:10 +0300] [13207] [INFO] Booting worker with pid: 13207
[2015-07-18 11:23:10 +0300] [13208] [INFO] Booting worker with pid: 13208


Дальше в браузере одновременно открываю 5 вкладок этого тестового сайта. Я ожидаю,
что эти вкладки в браузере отобразятся одновременно через 5 секунд, как задано в коде.
Но они появляются поочередно, через каждые 5 секунд. Как будто работая все равно в
один поток.

Что я делаю не так? :)
    


Ответы

Ответ 1



Я могу воспроизвести поведение в Firefox и Google Chrome. Браузер выполняет только один запрос по заданной ссылке за раз. Достаточно, использовать уникальные ссылки, чтобы увидеть, что gunicorn может обслуживать несколько запросов одновременно. Или вручную запустить несколько запросов (в этом случае не важно, одинаковые или разные ссылки). Вот Питон скрипт, который выполняет несколько одновременных http-запросов и в то же время открывает те же ссылки в браузере: #!/usr/bin/env python import webbrowser from multiprocessing.pool import ThreadPool try: from urllib2 import urlopen except ImportError: # Python 3 from urllib.request import urlopen urls = ['http://localhost:8000']*5 # same url urls += ['http://localhost:8000?unique=' + str(i) for i in range(5)] # uniq. urls pool = ThreadPool(len(urls) * 2) # make requests concurrently r = pool.map_async(lambda url: urlopen(url).read(), urls) pool.map(webbrowser.open_new_tab, urls) # open tabs in a browser r.get() Для тестирования можно использовать, простое wsgi-приложение: #file: wsgi_sleep.py import itertools import time def app(environ, start_response, ids=itertools.count(1)): status = '200 OK' headers = [('Content-type', 'text/plain')] data = "# request(s) per worker: " + str(next(ids)) headers.append(('Content-Length', str(len(data)))) start_response(status, headers) time.sleep(5) return [data] Его можно запустить как: $ gunicorn --threads 20 --access-logfile - wsgi_sleep:app Результат 127.0.0.1 - - [18/Jul/2015:19:15:09 +0000] "GET / HTTP/1.1" 200 26 "-" "Python-urllib/2.7" 127.0.0.1 - - [18/Jul/2015:19:15:09 +0000] "GET / HTTP/1.1" 200 26 "-" "Python-urllib/2.7" 127.0.0.1 - - [18/Jul/2015:19:15:09 +0000] "GET / HTTP/1.1" 200 26 "-" "Python-urllib/2.7" 127.0.0.1 - - [18/Jul/2015:19:15:09 +0000] "GET / HTTP/1.1" 200 26 "-" "Python-urllib/2.7" 127.0.0.1 - - [18/Jul/2015:19:15:09 +0000] "GET /?unique=0 HTTP/1.1" 200 26 "-" "Python-urllib/2.7" 127.0.0.1 - - [18/Jul/2015:19:15:09 +0000] "GET /?unique=1 HTTP/1.1" 200 26 "-" "Python-urllib/2.7" 127.0.0.1 - - [18/Jul/2015:19:15:09 +0000] "GET / HTTP/1.1" 200 26 "-" "Python-urllib/2.7" 127.0.0.1 - - [18/Jul/2015:19:15:09 +0000] "GET /?unique=3 HTTP/1.1" 200 27 "-" "Python-urllib/2.7" 127.0.0.1 - - [18/Jul/2015:19:15:09 +0000] "GET /?unique=2 HTTP/1.1" 200 26 "-" "Python-urllib/2.7" 127.0.0.1 - - [18/Jul/2015:19:15:09 +0000] "GET /?unique=4 HTTP/1.1" 200 26 "-" "Python-urllib/2.7" 127.0.0.1 - - [18/Jul/2015:19:15:14 +0000] "GET / HTTP/1.1" 200 27 "-" "Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:39.0) Gecko/20100101 Firefox/39.0" 127.0.0.1 - - [18/Jul/2015:19:15:14 +0000] "GET /?unique=2 HTTP/1.1" 200 27 "-" "Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:39.0) Gecko/20100101 Firefox/39.0" 127.0.0.1 - - [18/Jul/2015:19:15:14 +0000] "GET /?unique=4 HTTP/1.1" 200 27 "-" "Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:39.0) Gecko/20100101 Firefox/39.0" 127.0.0.1 - - [18/Jul/2015:19:15:14 +0000] "GET /?unique=1 HTTP/1.1" 200 27 "-" "Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:39.0) Gecko/20100101 Firefox/39.0" 127.0.0.1 - - [18/Jul/2015:19:15:14 +0000] "GET /?unique=3 HTTP/1.1" 200 27 "-" "Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:39.0) Gecko/20100101 Firefox/39.0" 127.0.0.1 - - [18/Jul/2015:19:15:14 +0000] "GET /?unique=0 HTTP/1.1" 200 27 "-" "Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:39.0) Gecko/20100101 Firefox/39.0" 127.0.0.1 - - [18/Jul/2015:19:15:19 +0000] "GET / HTTP/1.1" 200 27 "-" "Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:39.0) Gecko/20100101 Firefox/39.0" 127.0.0.1 - - [18/Jul/2015:19:15:24 +0000] "GET / HTTP/1.1" 200 27 "-" "Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:39.0) Gecko/20100101 Firefox/39.0" 127.0.0.1 - - [18/Jul/2015:19:15:29 +0000] "GET / HTTP/1.1" 200 27 "-" "Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:39.0) Gecko/20100101 Firefox/39.0" 127.0.0.1 - - [18/Jul/2015:19:15:34 +0000] "GET / HTTP/1.1" 200 27 "-" "Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:39.0) Gecko/20100101 Firefox/39.0" Лог показывает, что запросы от urllib клиента завершаются практически одновременно в то время как запросы от браузера для повторяющихся ссылок происходят последовательно (не связано с network.http.max-connections-per-server, возможно связано с настройками кэширования).

Ответ 2



Да, в общем-то, вы всё правильно делаете, вот только действительно работаете с одним потоком, так как ваш браузер уже имеет соединение с сервером. Попробуйте открыть разные браузеры.

четверг, 9 апреля 2020 г.

Android->Приложение->Потоки->AsyncStack Как сделать заставку для приложения,которая исчезнет через n секунд

#android #многопоточность #android_asynctask

                    
Подскажите,пожалуйста,как реализовать Asyntask для следующей задачи,у меня есть MainActivity,
в котором находится ListView, содержащий 1000 элементов. Мне нужно,чтобы перед MainActivity
запускалось другое Activity, содержащее только картинку, проходило 2 секунды, Activity
исчезает и появляется MainActivity. 
    


Ответы

Ответ 1



Это называется Splash Screen Вот пример реализации: public class SplashScreen extends Activity { private static int SPLASH_TIME_OUT = 2000; @Override protected void onCreate(Bundle savedInstanceState) { super.onCreate(savedInstanceState); setContentView(R.layout.activity_splash); new Handler().postDelayed(new Runnable() { @Override public void run() { Intent i = new Intent(SplashScreen.this, MainActivity.class); startActivity(i); finish(); } }, SPLASH_TIME_OUT); } } Смысл кода я думаю пояснять нет необходимости.

Что не так с определением дедлока?

#net #многопоточность #vbnet

                    
Код из вопроса Насколько случайны данные, сгенерированные таким образом? упал:



Т. е. в одном из методов произошёл вызов Release семафора, когда семафор был свободен.

Imports System.Threading

Module All
  Dim Value1 As Integer, Value2 As Integer
  Dim Sem1 As New SemaphoreSlim(1, 1), Sem2 As New SemaphoreSlim(1, 1)
  Dim Count As Integer

  Private Sub Inc(SemA As SemaphoreSlim, SemB As SemaphoreSlim, ByRef Value As Integer)
    Do
      Thread.Sleep(0)
      SemA.Wait()
      Interlocked.Increment(Count)
      SemB.Wait()
      Interlocked.Decrement(Count)
      Interlocked.Increment(Value)
      If SemB.CurrentCount = 0 Then SemB.Release() Else Interlocked.Increment(Count)
      If SemA.CurrentCount = 0 Then SemA.Release() Else Interlocked.Increment(Count)
    Loop
  End Sub

  Private Sub Init()
    Call (New Thread(Sub() Inc(Sem1, Sem2, Value1))).Start()
    Call (New Thread(Sub() Inc(Sem2, Sem1, Value2))).Start()
  End Sub

  Public Function GetRandBit() As Integer
    If Thread.VolatileRead(Count) = 2 Then
      Thread.VolatileWrite(Value1, 0)
      Thread.VolatileWrite(Value2, 0)
      Interlocked.Decrement(Count)
      Sem2.Release()
    End If

    Do Until Thread.VolatileRead(Count) = 2
      Thread.Sleep(16)
    Loop

    Return Thread.VolatileRead(Value1) And 1
  End Function

  Sub Main()
    Init()
    Do
      Console.Write(GetRandBit())
    Loop
  End Sub
End Module


Указанный в исключении метод _Lambda$__7-1 это метод второго потока:

[SpecialName]
internal void _Lambda\u0024__7\u002D1()
{
  All.Inc(All.Sem2, All.Sem1, ref All.Value2);
}


Из-за чего могла произойти такая ошибка?

Насколько я представляю, Sem2.Release() из GetRandBit вызывается только если оба
потока захватили по первой блокировке и никак не может возникнуть между проверкой на
количество блокировок и вызовом Release. На момент вызова Release одним из двух потоков,
он владеет обеими блокировками и никто посторонний семафоры не трогает. Что могло пойти
не так в этой схеме?
    


Ответы

Ответ 1



Похоже, всё-таки разобрался. Последовательность такая: Поток 2 владеет блокировкой 2 (1/2) Поток 1 владеет блокировкой 1 (1/2) Возникает дедлок (1/2) Управляющий поток снимает блокировку 2 (1/*) Поток 1 захватывает блокировку 2 (12/*) Поток 1 выполняет инкремент (12/*) Поток 1 снимает блокировку 2 (1/*) Поток 1 снимает блокировку 1 (0/*) Поток 2 захватывает блокировку 1 (0/*1) Поток 2 выполняет инкремент (0/*1) Поток 2 снимает блокировку 1 (0/*) Поток 1 захватывает блокировку 1 (1/*) Поток 1 захватывает блокировку 2 (12/*) Поток 2 проверяет наличие блокировки 2 (12/*) Блокировка есть, но поток 2 не знает, что она чужая Поток 1 делает инкремент (12/*) Поток 1 снимает блокировку 2 (1/*) Поток 2 пытается снять блокировку 2 и падает В скобках обозначено владение блокировками: 1 - захвачена блокировка 1 и поток находится внутри её блока 2 - захвачена блокировка 2 и поток находится внутри её блока * - блокировка 2 НЕ захвачена, но поток находится внутри её блока 0 - захваченных блокировок нет и поток не находится внутри какого-либо блока

четверг, 2 апреля 2020 г.

SDL_CreateThread и барьер памяти

#cpp #c #многопоточность #sdl2

                    
До запуска дополнительного потока создается семафор, который используется в создаваемом
потоке:

// Main thread:
SDL_sem* sem = SDL_CreateSemaphore(0);
SDL_CreateThread(...);

...

// Second thread:
SDL_SemWait(sem);


Вопрос - могут ли поменяться местами (компилятором/процессором) создание семафора
и потока? (и возникнет ошибка при обращении к ещё не созданному семафору в новом потоке).
    


Ответы

Ответ 1



Ключевое свойство всех оптимизаций компиляторов придерживающихся стандартов в том, что они не меняют результат вычислений, если в исходном коде нет UB. Исходя из этого можно быть вполне уверенным, что компилятор не переставит местами вызовы функций, покуда не будет точно уверен, что у них нет побочных эффектов (например, для gcc если они не объявлены с __attribute__(pure)). Таким образом, если вызов SDL_CreateSemaphore(0) идёт до SDL_CreateThread(...), то можно не боятся, что при исполнении они как-то поменяются... Конечно всегда может представить сферический компилятор в вакууме, который делает всё что ему вздумается, но это уже совсем другая история... Также может быть столь же сферическая библиотека реализующая свои вызовы, как отложенные, но это уже проблема библиотеки, и пользователь об этом беспокоиться не должен... И ни к любой вменяемой реализации SDL, ни к любому вменяемому компилятору это всё не относится...

Где использовать volatile

#java #многопоточность #concurrency

                    

Когда можно использовать volatile, если он не обеспечивает атомарность чтения и записи
как Atomic'и.
Мне не понятно в чем разница между volatile и synchronized, если и то и другое обеспечивает
синхронизацию между кэшами ядер процессора.

    


Ответы

Ответ 1



Модификатор volatile гарантирует видимость операций с полем и сохранение их последовательности. Волатильную переменную можно, например, использовать как флаг завершения работы потока: public class Main { private static volatile boolean run = true; public static void main(String[] args) throws Exception { new Thread(() -> { long x = 0; while (run) { System.out.println(++x); try { Thread.sleep(1000); } catch (InterruptedException exc) {} } }).start(); Scanner scanner = new Scanner(System.in); while (run) { String line = scanner.nextLine(); if ("exit".equals(line)) run = false; } } } Без модификатора volatile у вас нет гарантии, что поток выводящий значения переменной x когда-нибудь заметит, что главный поток изменил состояние переменной run. Синхронизация гарантирует видимость операций, сохранение их последовательности и атомарность. public class Main { private static int x = 0; private static int y = 1000; private static synchronized void transfer() { ++x; --y; } public static void main(String[] args) throws Exception { for (int a = 0; a < 10; a++) { new Thread(() -> { for (int b = 0; b < 100; b++) { transfer(); } }).start(); } System.out.println(x); System.out.println(y); } } Без модификатора synchronized у вас нет гарантии, что значение переменной x будет увеличено на столько же, на сколько уменьшено значение переменной y, так как порядок и продолжительность выполнения потоков непредсказуемы. Во втором примере volatile не поможет. В первом может помочь использование синхронизации, но тогда на каждой итерации один из потоков будет блокировать мьютекс, а второй потом будет останавливаться, пока мьютекс не будет освобождён. Это существенно медленнее проверки состояния волатильной переменной.

Ответ 2



Вам почти правильно ответили. Только оптимизации несколько другого рода. Дело в том, что переменная может быть закеширована для более быстрого доступа. Вариантов куча где именно. Например она может быть закеширована для каждого процессора отдельно, может быть закеширована в потоке, в принципе это зависит от того, какая реализация JVM используется. Volatile просто запрещает подобное кэширование. Теперь как его использовать. Volatile работает быстрее, потому что никаких блокировок не происходит. В этом его преимущество перед synchronized. Работа со ссылками атомарна, поэтому там в теории это может быть актуальным. А так да, конструкция относительно редкая, обычно при многопоточности используется что-то специфическое, дабы упростить работу с ней. UPD. Там "на заборах" пишут, что long и double которые не обязательно атомарные, становятся атомарными с volatile. Но официального подтверждения я этому не нашел, если у кого то есть точная информация по этому поводу, будет круто

вторник, 31 марта 2020 г.

Как из одного потока передать переменную в другой поток?

#java #многопоточность


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

public class s implements Runnable {
    public static void s(){
        System.out.println("тест");
    }
    @Override
    public void run() {
        s();
    }  
}

class Testpotok {
    public static void main(String[] args) {
        Thread t1 = new Thread(new s());
        t1.start();
        System.out.println("тест1");
    }
}

    


Ответы

Ответ 1



Swing, как и многие другие gui-библиотеки, однопоточен. При создании окна создаётся Event Dispatch Thread, внутри которого будет работать цикл событий и обработчики событий. Вы не должны пытаться из главного потока или любого другого потока взаимодействовать с элементами графического интерфейса - это приведёт к сбою. Вы не должны внутри обработчиков событий запускать потоки - это приведёт к сбою. Если вам надо из другого потока изменить, например, текст метки, то придётся создать задание для EDT: SwingUtilities.invokeLater(() -> someLabel.setText("Hello")); Если вам нужно внутри обработчика нажатия на кнопку запустить на выполнение длительную задачу, придётся использовать SwingWorker: SwingWorker worker = new SwingWorker() { @Override protected void doInBackground() throws Exception { // Выполняется в отдельном потоке } @Override protected void done() { // Выполняется в EDT после завершения doInBackground } }; worker.execute();

воскресенье, 29 марта 2020 г.

Синхронизация вывода потоков POSIX

#cpp #c #многопоточность #pthread #posix


Нужно что бы два потока параллельно печатали на экран. (Первый поток печатает числа
1,2,3...10 Второй - 100,200,300...1000). Причём вывод должен быть синхронизирован:
сначала родительский поток выводит первую строку, затем дочерний первую, затем родительский
вторую строку, затем дочерний вторую и т.д.(100,1,200,2,300,3...) Использовать нужно
мьютексы.

pthread_mutex_t mut;

void printt(int i){
   pthread_mutex_lock(&mut);
   cout<


Ответы

Ответ 1



Предупреждение: согласно POSIX данное решение даёт UB [1], хотя и, судя по всему, он работает в реализации pthreads от glibc/linux; он приведен лишь для справки/как идея и не должен использоваться. Спасибо @VTT за замечание. Для решения задачи нужно количество мьютексов равное количеству потоков. Идея в том, чтобы мьютексы захватывались разными потоками попеременно, т.е. входя в критические секции нужно захватывать свой мьютекс, а при выходе отпускать мьютекс следующего потока. Само собой, перед началом выполнения свободным должен быть отпущен только один мьютекс. В итоге получается нечто следующее: pthread_mutex_t mut1; pthread_mutex_t mut2; void printt1(int i){ pthread_mutex_lock(&mut1); cout<

Ответ 2



Функция printt печатает локальную переменную i, к которой у других потоков нет доступа, а синхронизация не ожидает завершения записи предыдущей строки другим потоком. Также полностью отсутствует обработка ошибок.

Ответ 3



Без сигналов , чтоб другой поток проснулся неудобно, придумал только сон. Вроде бы пашет. Нужна переменная для знака кому какая очередь. // g++ -pthread mutex-semaph.cpp #include # include # include pthread_mutex_t mut; int queue = 1 ; // или 2 void printt(int i, int q){ Again : pthread_mutex_lock(&mut); if(queue == q) std::cout<

Ответ 4



Сделайте разделяемую volatile переменную и mutex. Присвойте переменной 1, что означает печать будет проводить первый поток. Запустите потоки. В каждом потоке в цикле захватываете mutex и читаете значение переменной. Если значение переменной в первом потоке равно 1, то он печатает данные и присваивает переменной 2. Аналогично, второй поток производит печать если значение переменной равно 2, после чего устанавливает ее в 1. (Обратите внимание, чтение и модификация переменной защищены mutex-ом.) Затем поток в любом случае снимает блокировку и (для оптимизации эффективности) вызывает pthread_yield. Конкретно в вашем случае с внешним циклом и анализом очередности печати внутри printt() (очевидно, с целью не вытаскивать работу с mutex на уровень управления циклом) в этой функции нужно дождаться, пока общая переменная не примет нужного значения. Впрочем, хватит общих слов. Вот немного модифицированный ваш код из текста вопроса. avp@avp-ubu1:hashcode$ cat t-seq-pri.c #ifdef __cplusplus #include using namespace std; #else #define _GNU_SOURCE #include #endif #include pthread_mutex_t mut; volatile int turn; void printt(int i, int q){ int pri = 0; do { pthread_mutex_lock(&mut); if (turn == q) { #ifdef __cplusplus cout<

Java Где хранится volatile переменная

#java #многопоточность #volatile


Всегда думал что volatile переменные в Java хранятся в MetaSpace, недавно на собеседовании
мне сказали что это неверно. Так вот вопрос: где они хранятся?
    


Ответы

Ответ 1



Даже интересно, откуда у вас могла возникнуть такая мысль. В метаспэйсе, как и следует из его названия, хранятся описания типов, а не данные. За исключением разве что констант. Данные хранятся либо в стеке, либо в куче. Изредка в нативной памяти. Так как модификатор volatile может применяться только к полям, то волатильные значения всегда будут в куче.

воскресенье, 15 марта 2020 г.

“Тонкая блокировка” в хэш-таблицах

#алгоритм #многопоточность


В чем суть "тонкой блокировки" в хэш-таблицах и как её реализовать?
    


Ответы

Ответ 1



Суть в том, что с каждым входом (цепочкой синонимов хэш-функции) связывают свой mutex. Реализация на Си тривиальна: struct hash_entry { pthread_mutex_t mutex; struct hash_item *list; // если, например, используете односвязный список }; При создании таблицы struct hash_entry table[HASH_SIZE]; устанавливаете table[i].list в NULL и вызываете pthread_mutex_init(&table[i].mutex, 0); (если нет особых требований к атрибутам мьютекса). Update После комментария от @VladD решил добавить немного рассуждений на тему параллельной работы с хэш-таблицей. В принципе для действий, затрагивающих всю таблицу, можно добавить rwlock (см., например, man pthread_rwlock_init и др. man для pthread_rwlock...). Тогда функции, требующие "тонкой блокировки" будут вызывать pthread_rwlock_t tablock; ... pthread_rwlock_rdlock(&tablock); // совместная блокировка таблицы int i = hash_func(key, keylen) % HASH_SIZE; pthread_mutex_lock(&table[i].mutex); // блокировка цепочки // действия с цепочкой синонимов ... pthread_mutex_unlock(&table[i].mutex); pthread_rwlock_unlock(&tablock); а функции, которые затрагивают всю таблицу pthread_rwlock_wrlock(&tablock); // монопольная блокировка таблицы // действия с таблицей ... pthread_rwlock_unlock(&tablock); т.е. получаем схему писатель-читатели с вложенной монопольной блокировкой цепочек коллизий. Однако, на мой взгляд, в большинстве случаев подобные схемы (и вообще "тонкая блокировка") не оправдывают возлагаемых на них надежд из-за накладных расходов на блокировку. Реально, если не удается вообще уйти (а к этому всегда надо стремиться (KISS-принцип)) от многопоточной работы с таблицей, выгодней просто блокировать ее всю, независимо от вида операций, поскольку обычно цепочки синонимов короткие и большинство операций с элементами таблицы весьма быстрые. Еще одним доводом в пользу такого подхода может быть тривиальная экономия памяти (размер объекта pthread_mutex_t 40 байт на 64-bit X_86), которая (как это ни странно, на первый взгляд) часто приводит к ускорению программы (поскольку доступ к памяти может быть раз в 100 медленней доступа к кэшу). Если возникают сомнения по поводу времени работы самой хэш-функции от ключа (например, известно, что ключи на самом деле весьма длинные и используется какая-то изощренная схема хэширования), то ее вычисление можно проводить до блокировки таблицы (а вот вычисление самого индекса, конечно же, уже получив блокировку).

пятница, 13 марта 2020 г.

Самописный СМТП-сервер

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


Я написал СМТП сервер, но работает он некорректно вот в каком моменте: 
если я 1 раз подключусь к нему через телнет, он ответит нормально (рис 1-2.), но
если я еще раз подключусь через другую консоль cmd, то я появится только черный экран
и все(рис3):


(рис 1)


(рис 2)


(рис 3)

В последнем случае что-бы я ни вводил, на экране ничего не отобразится, и даже не
залогируется моим СМТП

В первом же случае, мои команды отображаются в консоли, они логируюстя сервером и
я вижу в этой консоли ответ (рис 4):


(рис 4)

Т.е., что я вижу, один клиент к нему подключается, работает нормально сервер, во
всех остальных случаях он работает не корректно (верно?)

Вот как я реализовал СМТП-сервер (это win-сервис):

 protected override void OnStart(string[] args)
        {
SmtpHelper s = new SmtpHelper(this);
System.Threading.ThreadPool.QueueUserWorkItem(new System.Threading.WaitCallback(StartListen),
(object)s );
}

 void StartListen(object s)
        {
            try
            {
                var a = (SmtpHelper)s;
                a.Listen(); //запуск 
            }
            catch (Exception ex)
            {
                l.Write("Error (StartListen(object s)): " + ex.ToString());
                throw;
            }

        }




public void Listen()
        {
            try
            {
                SMTP_Listener = new TcpListener(IPAddress.Any, port); 
                SMTP_Listener.Start();

                while (true)
                {
                    clientSocket = SMTP_Listener.AcceptSocket();

                    _sessionId = clientSocket.GetHashCode().ToString();

                    _email.sessionId = Convert.ToInt32(_sessionId);


                    StartProcessing(newController);
                    l.Write("we are there");

                   // System.Threading.ThreadPool.QueueUserWorkItem(new System.Threading.WaitCallback(ClientThread),
SMTP_Listener.AcceptTcpClient());
                }

            }
            catch (Exception ex)
            {
                l.Write("SMTP Listen Error: " + ex.ToString());
                throw;
            }
        }



void StartProcessing(ClientSessionController newController)
        {

            try
            {
                m_ConnectedIp = ParseIP_from_EndPoint(clientSocket.RemoteEndPoint.ToString());
                m_ConnectedHostName = GetHostName(m_ConnectedIp);

                _email.ip = m_ConnectedIp;
                _email.port = 25;


                SendData("220 " + System.Net.Dns.GetHostName() + " Service ready\r\n");

                //if (!clientSocket.Connected)
                //    clientSocket.Connect(IPAddress.Any, port);

                //РАБОТА С ВХОДНЫМИ ДАННЫМИ
                while (true)
                {
                    //если есть данные, то считаем их
                    if (clientSocket.Available > 0)
                    {
                        //получение команды от клиента
                        string lastCmd = ReadLine();

                        //парсим команду от клиента (HELO, RCPT, DATA etc.)
                        if (lastCmd.Trim() != String.Empty)
                            ProceedCommand(lastCmd, newController);
                        //break;
                    }
                    else
                    {
                      //dump:  l.Write("[Socket isn't available now]");
                    }
                }               
            }
            catch (Exception ex)
            {
                throw;
            }

        }


ВОПРОС: 

Подскажите пожалуйста, как мне сделать так, чтобы мой сервис мог одинаково обслужить
всех подключаемых клиентов, а не только первого?
    


Ответы

Ответ 1



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

Как отслеживать запуски GC и учитывать их при логировании из разных потоков?

#c_sharp #net #многопоточность #логирование


Из разных потоков вызывается метод и надо логировать его работу. Т.к. GC при сборке
мусора в WinForms и WPF приложениях может приостановить работу потоков приложения,
то это надо учитываться при логировании.
Как можно отслеживать запуски GC?
    


Ответы

Ответ 1



Для того чтобы отслеживать GC надо вызвать метод RegisterForFullGCNotification, а также WaitForFullGCApproach и WaitForFullGCComplete. using System.Diagnostics; using System.Collections.Concurrent; using System.Threading; using System.Runtime.CompilerServices; class LogLine { public int Number; public int ThreadId; public long Ticks; public object Value; public long GCMemory; public int[] GCCollections; public LogLine(int num, object value, long ticks) { this.Number = num; this.Value = value; this.Ticks = ticks; this.ThreadId = Environment.CurrentManagedThreadId; this.GCMemory = GC.GetTotalMemory(false); int[] arr = new int[GC.MaxGeneration + 1]; for (var i = 0; i < arr.Length; i++) arr[i] = GC.CollectionCount(i); this.GCCollections = arr; } } class Log : IDisposable { Stopwatch sw = Stopwatch.StartNew(); BlockingCollection lines = new BlockingCollection(); void Add(LogLine line) { if (!lines.IsAddingCompleted) lines.Add(line); } public Log() { GC.RegisterForFullGCNotification(1, 1); new Thread(() => { while (!lines.IsAddingCompleted) { Add(new LogLine(-1, GC.WaitForFullGCApproach(), sw.ElapsedTicks)); Add(new LogLine(-2, GC.WaitForFullGCComplete(), sw.ElapsedTicks)); } }).Start(); } public void WriteLine(object value = null, [CallerLineNumber] int cnumber = 0) { Add(new LogLine(cnumber, value, sw.ElapsedTicks)); } void IDisposable.Dispose() { GC.CancelFullGCNotification(); Add(new LogLine(-3, "Disposed", sw.ElapsedTicks)); lines.CompleteAdding(); } public IEnumerable ToCsv() { var s = ",\t "; yield return String.Concat( "Number", s, "Ticks", s, "ThreadId", s, "GCMemory", s, "GCCollections", s, "Value"); foreach (var l in lines) yield return String.Concat( l.Number, s, l.Ticks, s, l.ThreadId, s, l.GCMemory, s, String.Join(";", l.GCCollections), s, l.Value); } } Для теста в двух потоках создаем и заполняем список массивов по 100 тыс. int var log = new Log(); using (log) { Parallel.For(0, 2, i => { var lst = new List(); log.WriteLine("new List"); try { while (true) lst.Add(new int[100000]); } catch (OutOfMemoryException) { log.WriteLine("OutOfMemory; lst.Count=" + lst.Count); } }); } Выводим собранные данные foreach (var line in log.ToCsv()) Console.WriteLine(line); Результат (получен в C# Interactive, Microsoft (R) Roslyn C# Compiler version 1.1.0.51204) Number, Ticks, ThreadId, GCMemory, GCCollections, Value 55, 4030, 6, 6905852, 1;0;0, new List 55, 4426, 9, 6905852, 1;0;0, new List -1, 274370, 11, 1481676908, 6;5;5, Succeeded -2, 274430, 11, 1482877004, 6;5;5, Succeeded -1, 277392, 11, 1528080620, 6;5;5, Succeeded 58, 355462, 6, 1528087152, 8;7;7, OutOfMemory; lst.Count=1888 58, 355462, 9, 1528087152, 8;7;7, OutOfMemory; lst.Count=1920 -2, 317257, 11, 1528087152, 8;7;7, Succeeded -3, 356158, 6, 1528095344, 8;7;7, Disposed

Поток потребляет очень много ресурсов C#

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


Мне необходимо получать в фоне какая раскладка клавиатуры сейчас активна в системе.

Само получение раскладки реализовано так:

[DllImport("user32.dll", SetLastError = true)]
static extern int GetWindowThreadProcessId(
   [In] IntPtr hWnd,
   [Out, Optional] IntPtr lpdwProcessId
);

[DllImport("user32.dll", SetLastError = true)]
static extern IntPtr GetForegroundWindow();

[DllImport("user32.dll", SetLastError = true)]
static extern ushort GetKeyboardLayout(
   [In] int idThread
);


Непосредственно функция:

ushort GetKeyboardLayout()
{
    return GetKeyboardLayout(GetWindowThreadProcessId(GetForegroundWindow(), IntPtr.Zero));
}


Я создал поток такого вида:

while (true)
{
    if (GetKeyboardLayout() == 1033)
    {
        this.BackgroundImage = bt1;
        this.notifyIcon1.Icon = WindowsFormsApplication1.Properties.Resources.engico;
    }
    else
    {
        this.BackgroundImage = bt2;
        this.notifyIcon1.Icon = WindowsFormsApplication1.Properties.Resources.ruico;
    }
}


Собственно запуск потока:

Thread thread1 = new Thread(tickness);
thread1.IsBackground = true;
thread1.Priority = ThreadPriority.Lowest;
thread1.Start();


Однако, такая реализация потребляет около 15% CPU и в целом "кушает" очень много.

Подскажите, пожалуйста, более экономичный вариант работы потока.
    


Ответы

Ответ 1



Все правильно, у вас поток постоянно молотит в бесконечном цикле и это отъедает изрядную часть процессорного времени. Таймер и вызов вашего кода, допустим, каждые 100мс определенно понизят нагрузку. Вам нужно решить каков допустимый интервал между сменой раскладки и реакцией программы и использовать его. Без таймера можно на WaitForSingleObject в цикле, если я ничего не путаю. Но это тоже просто задержка, которая не жрет процессор.

Многопоточное подключение к SSH

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


Наставьте на путь истинный, подскажите, что можно сделать/почитать по этому поводу?

Есть список объектов типа IP.

class IP
{
    public string host;
    public string login;
    public string password;
}


И метод, который перебирает этот список в параллельном цикле, в котором происходит
подключение к SSH. 

public void Cycle()
{
  int port = 22;

  Parallel.ForEach(ipList, ipObject =>
  {
       string ip = ipObject.host;
       string login = ipObject.login;
       string password = ipObject.password;

       SshConnect(ip, port, login, password);
  });
}


Метод SSHConnect реализован с помощью библиотеки Renci.SSHNet. Выглядит следующим
образом:

private bool SshConnect(string host, int port, string login, string password)
 {
     bool flag = true;
     try
     {
          var client = new SshClient(host, port, login, password);
          client.Connect();
          client.Disconnect();
     }
     catch
     {
          flag = false;
     }
     return flag;
 }


Все бы ничего, но скорость перебора и подключения к SSH ужасно медленная. Что можно
сделать, чтобы увеличить скорость?

P.S. Раньше у меня был алгоритм, который работал, как минимум, в 2 раза быстрее этого.

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

На другом форуме сказали, что так делать нежелательно и немного почитав, я решил
попробовать Paralel.ForEach, но результат не оправдал ожиданий.
    


Ответы

Ответ 1



Основное время при подключении по ssh тратит не процессор, а сеть - пакет TCP SYN слишком долго идет до сервера и обратно. В таких условиях надо и правда создавать как можно больше потоков, Paralel.ForEach же ограничивает их число. У вас есть список ssh-серверов, с которыми надо что-то сделать - или вы сканируете сеть? Во втором случае можно сделать оптимизацию. Сначала можно проверить, открыт ли порт 22 - а уже потом подключаться по ssh там, где он открыт. Это можно сделать параллельно, но вовсе без потоков, при помощи асинхронного программирования: var taskList = ipList.Select(async ipObject => { try { using (var c = new TcpClient()) await c.ConnectAsync(ipObject.host, 22); retrun ipObject; } catch { return null; } }).ToArray(); Task.WaitAll(taskList); ipList = taskList.Select(task => task.Result).Where(ipObject => ipObject != null).ToList(); // Теперь в ipList остались только хосты, которые отвечают на порту 22

воскресенье, 8 марта 2020 г.

Создание анимации через потоки на Android

#java #android #многопоточность #canvas #анимация


Здравствуйте, стоит задача, на языке Java под ОС Android нужно написать задачу: создать
снеговика и сделать так, чтобы его составляющие объекты меняли цвет с разной скоростью.
Смена цветов должна быть реализована в потоке.
У меня получилось написать код, но не понимаю, как именно эта анимация должна в этот
поток помещаться. Вот собственно полный код, написан в Eclipse:

package ru.ucheba;

import android.app.Activity;
import android.content.Context;
import android.graphics.Canvas;
import android.graphics.Paint;
import android.graphics.Paint.Style;
import android.os.Bundle;
import android.os.SystemClock;
import android.view.View;

public class Zadanie2 extends Activity {
    int a = 255;
    int r = 100;
    int g = 0;
    int b = 0;

    @Override
    protected void onCreate(Bundle savedInstanceState) {
        super.onCreate(savedInstanceState);
        setContentView(new Panel(this));

    }

    class Panel extends View {
        public Panel(Context context) {
            super(context);
        }

        @Override
        public void onDraw(Canvas canvas) {
            super.onDraw(canvas);
            Paint p1 = new Paint();
            p1.setStyle(Style.FILL);

            p1.setARGB(a, r, g, b);
            canvas.drawCircle(270, 170, 70, p1);
            for (int i = 0; i < 5; i++) {
                r--;
                SystemClock.sleep(10);
                invalidate();
            }

            for (int j = 0; j < 15; j++) {
                p1.setARGB(a, r, g, b);
                g--;
                canvas.drawCircle(270, 310, 100, p1);
                SystemClock.sleep(20);
                invalidate();
            }

            for (int k = 0; k < 30; k++) {
                p1.setARGB(a, r, g, b);
                b--;
                canvas.drawCircle(270, 510, 150, p1);
                SystemClock.sleep(40);
                invalidate();
            }
        }

        class Task extends Thread {
            @Override
            public void run() {

            }
        }
    }
}


Чем заполнять метод run() понимаю, но ругается либо на canvas, p1, или invalidate().
Прошу помощи
    


Ответы

Ответ 1



Ваш код принципиально неверен. В методе onDraw() нельзя вызывать sleep() и invalidate(). Это все вместе с изменением переменных argb должно происходить в отдельном потоке, а метод onDraw() только отрисовывать экран в соответствии с их состоянием. @Override public void onDraw(Canvas canvas) { super.onDraw(canvas); p.setARGB(a, r1, g1, b1); canvas.drawCircle(270, 170, 70, p); p.setARGB(a, r2, g2, b2); canvas.drawCircle(270, 310, 100, p); p.setARGB(a, r3, g3, b3); canvas.drawCircle(270, 510, 150, p); } Перенести в конструктор и сделать переменной класса: p = new Paint(); p.setStyle(Style.FILL); В run() меняйте rX, gX, bX как вам надо и вызывайте синхронно с основным потоком invalidate() когда надо перерисовать.

Ответ 2



Выполнил вот таким образом, может кому-то поможет. Думаю, этот код можно написать гораздо лучше. Поэтому кто сможет, укажите на ошибки, буду благодарен. import android.app.Activity; import android.content.Context; import android.graphics.Canvas; import android.graphics.Color; import android.graphics.Paint; import android.os.Bundle; import android.view.View; public class Zadanie2 extends Activity { int[] Colo = { 10, 20, 30 }; // массив со значениями исп. для цветов // снеговика public Thread myThread1, myThread2, myThread3; // создание переменных для потоков @Override protected void onCreate(Bundle savedInstanceState) { super.onCreate(savedInstanceState); setContentView(new Panel(this)); // использование класса Panel в качестве активити myThread1 = new Thread(new Runnable() { @Override public void run() { int znak = 1; while (true) { if ((znak > 0) && (Colo[0] > 250)) { znak = -znak; } if ((znak < 0) && (Colo[0] < 5)) { znak = -znak; } Colo[0] += znak; try { Thread.sleep(30); } catch (InterruptedException e) { e.printStackTrace(); } } } }); myThread2 = new Thread(new Runnable() { @Override public void run() { int znak2 = 1; while (true) { if ((znak2 > 0) && (Colo[1] > 250)) { znak2 = -znak2; } if ((znak2 < 0) && (Colo[1] < 5)) { znak2 = -znak2; } Colo[1] += znak2; try { Thread.sleep(30); } catch (InterruptedException e) { e.printStackTrace(); } } } }); myThread3 = new Thread(new Runnable() { @Override public void run() { int znak3 = 1; while (true) { if ((znak3 > 0) && (Colo[2] > 250)) { znak3 = -znak3; } if ((znak3 < 0) && (Colo[2] < 5)) { znak3 = -znak3; } Colo[2] += znak3; try { Thread.sleep(30); } catch (InterruptedException e) { e.printStackTrace(); } } } }); // запуск потоков myThread1.start(); myThread2.start(); myThread3.start(); } class Panel extends View { public Panel(Context context) { super(context); } @Override public void onDraw(Canvas canvas) { super.onDraw(canvas); float w, h, cx, cy, radius; // переменные, для адаптивного расположения снеговика w = getWidth(); // считывает ширину h = getHeight(); // считывает высоту cx = w / 2; cy = h / 2; // для ориентации экрана if (w > h) { radius = h / 8; } else { radius = w / 8; } Paint p1 = new Paint(); p1.setStyle(Paint.Style.FILL); p1.setColor(Color.rgb(Colo[0], 255, 255)); canvas.drawCircle(cx, cy - h / 3, radius, p1); Paint p2 = new Paint(); p2.setStyle(Paint.Style.FILL); p2.setColor(Color.rgb(255, Colo[1] * 2, 255)); canvas.drawCircle(cx, cy - h / 3 + radius * 2, (float) (radius * 1.5), p2); Paint p3 = new Paint(); p3.setStyle(Paint.Style.FILL); p3.setColor(Color.rgb(255, 255, Colo[2] * 3)); canvas.drawCircle(cx, cy - h / 3 + radius * 5, radius * 2, p3); invalidate(); // перерисовка объектов } } }

Как усыпить текущий поток с возможностью возобновления из другого потока в Qt?

#cpp #qt #многопоточность #qt_faq


Имеется поток, который иногда засыпает(QThread::sleep) на довольно продолжительное
время (около 5 секунд). Если пользователь закрывает программу, то основной поток пытается
завершить уснувший, и ожидает его действительного завершения с помощью функции QThread::wait.
Но 5 секунд - достаточно продолжительное время, чтобы пользователь занервничал. В связи
с этим вопрос: как разбудить поток, который внутри себя вызвал функцию QThread::sleep?
    


Ответы

Ответ 1



Qt не предоставляет средств для пробуждения спящего потока. Но вместо этого можно использовать класс QWaitCondition с некоторой доработкой. Напишем для этого специальный класс. WakeableSleep.h: #ifndef WAKEABLESLEEP_H #define WAKEABLESLEEP_H #include #include #include #include /** * @brief Класс, который позволяет временно усыпить поток с возможностью пробуждения из другого потока. * * Класс можно создать в любом потоке. При вызове метода \ref sleep поток приостанавливается * на время, переданное с параметром. При вызове метода \ref wake из другого потока целевой * поток возобновляет выполнение независимо от истекшего времени. * \threadsafe */ class WakeableSleep : public QObject { Q_OBJECT public: explicit WakeableSleep(QObject *parent = 0); /** * @brief Усыпить текущий поток на milleseconds миллисекунд. * @param milliseconds Время сна. */ void sleep(quint32 milliseconds); /** * @brief wake Пробудить целевой поток из другого потока. */ void wake(); private: QMutex mutex; QWaitCondition waitCondition; }; #endif // WAKEABLESLEEP_H WakeableSleep.cpp: #include "wakeablesleep.h" WakeableSleep::WakeableSleep(QObject *parent) : QObject(parent){} void WakeableSleep::sleep(quint32 milliseconds) { mutex.lock(); waitCondition.wait(&mutex, milliseconds); mutex.unlock(); } void WakeableSleep::wake() { mutex.lock(); waitCondition.wakeAll(); mutex.unlock(); } Теперь вместо метода QThread::sleep можно использовать методы этого класса следующим образом: WakeableSleep sleeper; void Thread1() { sleeper.wake(); } void Thread2() { sleeper.sleep(5000); }

Realm and Threading

#java #android #многопоточность #realm #realm_android


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

public class QuoteText extends RealmObject {
    private long id;
    private String quoteText;
    ...
}


Для изолирования слоя работы с этой бд я решила применить обобщенный вариант паттерна
Repository. Написала интерфейс:

public interface QuoteRepository {
    List getListOfQuoteText();
}


И класс, реализующий этот интерфейс:

public class QuoteDataRepository implements QuoteRepository {
    private final Realm realm;  
    @Override
    public List getListOfQuoteText() {
        return realm.where(QuoteText.class).findAll();
    }
}


Соответственно во фрагменте для получения списка QuoteText из бд:

QuoteDataRepository quoteDataRepository = new QuoteDataRepository();
List quoteTexts = quoteDataRepository.getListOfQuoteText();


Все бы ничего, но хотелось бы все эти запросы делать не в UI потоке. Как засунуть
это все в другой поток? (Особенно интересно: можно ли совместить способы асинхронных
запросов предлагаемые Realm (назнчаение слушателей, или запрос с использованием onSuccess(),
onError() и прочее) и изоляцию слоя работы с бд). 

Спасибо за помощь!

Правка: важен момент именно абстракции кода работы с бд и реализации асинхронных
запросов.
    


Ответы

Ответ 1



Красивее всего получится если вы будете соединять Realm + RxJava + Retrolambda. Например вот такое может получиться: public Observable> getFeedVkPostsSortedAsync(String field, Sort order) { return mRealm.where(VkPost.class) .equalTo(VkPost.FIELD_IS_IN_FEED, true) .equalTo(VkPost.FIELD_IS_IN_FAVORITES, true) .findAllSortedAsync(field, order) .asObservable() .filter(RealmResults::isLoaded) .filter(RealmResults::isValid) //опционально и на всякий случай отвязываем объекты от реалма и //возвращаем обычные объекты в обычном листе .flatMap(realmResults -> Observable.just(mRealm.copyFromRealm(realmResults))); } Данный Observable, запущенный из основного потока выполнит запрос в БД асинхронно, отсеет объекты не удовлетворяющие двум условиям значений полей объектов (VkPost.FIELD_IS_IN_FEED, VkPost.FIELD_IS_IN_FAVORITES) и проверит что объекты возвращаемые полностью готовы для использования. И будет эмитить новую выборку при каждом изменении в ней.

Ответ 2



У Realm есть возможность создавать асинхронные запросы. Для этого нужно вместо findAll() вызывать findAllAsync(). Поправлю Ваш интерфейс, потому что оба этих метода возвращают не List, а RealmResults - это объект стандарта Future: public interface QuoteRepository { RealmResults getListOfQuoteTextAsync(); } Для того, чтобы получить уведомление об окончании загрузки нужно добавить подписку на экземпляр RealmResults: RealmResults result = quoteDataRepository.getListOfQuoteTextAsync(); result.addChangeListener(new RealmChangeListener() { @Override public void onChange(RealmResults results) { // метод будет вызван, когда запрос будет выполнен или при обновлении данных } }); Помимо этого, можно убедиться в завершении загрузки вызвав метод isLoaded(): if (result.isLoaded()) { // данные загружены } Получение результата Для работы с результатами запросов (в том числе асинхронными) Realm предоставляет специализированные адаптеры, которые нужно добавить в зависимости в build.gradle: dependencies { compile 'io.realm:android-adapters:1.4.0' } После этого нужно создать наследника от RealmRecyclerViewAdapter, который будет работать с Вашим ViewHolder. Продемонстрирую использование на примере из документации public class MyFragment extends Fragment { private Realm realm; private RecyclerView recyclerView; @Override public View onCreateView(LayoutInflater inflater, ViewGroup container, Bundle savedInstanceState) { realm = Realm.getDefaultInstance(); View root = inflater.inflate(R.layout.fragment_view, container, false); recyclerView = (RecyclerView) root.findViewById(R.id.recycler_view); // установка Вашего адаптера для RecyclerView recyclerView.setAdapter(new MyRecyclerViewAdapter(getActivity(), // установка результата асинхронного запроса quoteDataRepository.getListOfQuoteTextAsync())); // ... return root; } @Override public void onDestroyView() { super.onDestroyView(); realm.close(); } } В общем случае для получения (и отображения) результата внутри Fragment/Activity нужно осуществить подписку на RealmResults. Сделать это нужно именно в вызывающем коде для возможности отписаться от уведомлений (так как никто кроме вызывающего кода не знает, когда запрос для него уже неактуален). Этот вариант будет выглядеть так: public class MyFragment extends Fragment { private Realm realm; private RealmResults results; @Override public View onCreateView(LayoutInflater inflater, ViewGroup container, Bundle savedInstanceState) { realm = Realm.getDefaultInstance(); View root = inflater.inflate(R.layout.fragment_view, container, false); results = quoteDataRepository.getListOfQuoteTextAsync(); // добавление подписки на получение результата results.addChangeListener(new RealmChangeListener() { @Override public void onChange(RealmResults results) { // метод будет вызван, когда запрос будет выполнен или при обновлении данных // здесь можно обновлять экран или делать другой полезный код } }); return root; } @Override public void onDestroyView() { super.onDestroyView(); // при уничтожении фрагмента нужно отписаться от уведомлений results.removeChangeListeners(); realm.close(); } }

Ответ 3



Вот эта статья(-и) дала(-и) ответы на все мои вопросы! Спасибо автору! https://medium.com/@Viraj.Tank/realm-integration-in-android-best-practices-449919d25f2f#.9735g4ojc

суббота, 7 марта 2020 г.

Асинхронная загрузка изображений без остановки ui wpf

#c_sharp #wpf #многопоточность #binding #async


Есть стэк, который биндится к списку изображений. 


    
        
            
                
            
        
        
            
                
                    
                
            
        
    



Как мне асинхронно подгружать в него изображения (Большой объем), чтобы ui не останавливался?
Пробовал в коллекцию через асинхронный метод добавлять, но поток все равно останавливается.
Во view добавлял объекту IsAsync, тоже не помогает.

private async void GenerateTop()
{
    Screenshots = new ObservableCollection();
    await Task.Run(() =>
    {
        LastAdded.Add(...)...
    }
}

    


Ответы

Ответ 1



Вы должны разделить модель и представление. Функции загрузки модели должны бежать в неосновном потоке, и там пусть загружают что угодно как угодно медленно. Представление (в случае, если вы используете MVVM, это ваша VM) получает от модели любым способом извещение о том, что данные загрузка произошла и есть новые данные для показа, и (в главном потоке!) обновляет VM-список изображений. Таким образом ваш UI не будет подвисать, ваш код будет будет асинхронным, а волосы — мягкими и шелковистыми. Если подвисает загрузка картинки, попробуйте и правда использовать IsAsync в привязке. Единственная тонкость — у вас сейчас IsAsync грузит асинхронно строку, а конверсия в картинку выполняется в UI-потоке. Попробуйте указать конвертер: class ImageSourceLoadingConverter : IValueConverter { public object Convert(object value, Type targetType, object p, CultureInfo ci) => new BitmapImage(new Uri((string)value)); public object ConvertBack(object value, Type targetType, object p, CultureInfo ci) => throw new NotSupportedException(); } в этом случае, кажется, конверсия будет производиться асинхронно.

Ответ 2



Сделайте обычный метод(не асинхронный). Заполните в нем всю коллекцию и попробуйте так: далее: private void GenerateTop() { Screenshots = new ObservableCollection(); LastAdded.Add(...)... }

четверг, 5 марта 2020 г.

C++ std::thread. Как обратить и завершить поток

#cpp #многопоточность


Всем привет, 

Как обратиться к потоку по id и попросить его "завершиться"?

Использую стандарт С++11 и библиотеку  

Спасибо)
    


Ответы

Ответ 1



Никак. std::thread такое не предусматривает. Используйте atomic, [shared_]future::wait_for или другие примитивы чтобы просигнализировать коду внутри потока о том что он должен завершиться.

среда, 4 марта 2020 г.

Параллельное выполнение в случае вложенных циклов

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


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

double optResult = 0;
double result = 0;
for (int i = 5; i < 50; i = i + 5)
{
    for (int j = 5; j < 50; j = j + 5)
    {
        if (i < j)
        {
            optResult = OptimStart(i,j);
            if (optResult > result)
            {
                result = optResult;
            }
        }

    }
}

    


Ответы

Ответ 1



Я бы сделал как то так void Main() { var tuples = new List>(); for (int i = 5; i < 50; i = i + 5) { for (int j = i+5; j < 50; j = j + 5) { tuples.Add(Tuple.Create(i, j)); } } var max = tuples.AsParallel().Max(t=>Foo(t.Item1, t.Item2)); Console.WriteLine(max); } int Foo(int i, int j) { Thread.Sleep(1000); return i+j; } Если комбинаций слишком много и не хочется держать их в памяти, то void Main() { var max = GetItems().AsParallel().Max(t=>Foo(t.Item1, t.Item2)); Console.WriteLine(max); } int Foo(int i, int j) { Thread.Sleep(1000); return i+j; } IEnumerable> GetItems() { for (int i = 5; i < 50; i = i + 5) { for (int j = i + 5; j < 50; j = j + 5) { yield return Tuple.Create(i, j); } } }

Ответ 2



Я бы единственное предложил бы вложенный цикл начинать не с 5, а с i + 5, и тогда внутреннее условие if (i < j) можно убрать: double optResult = 0; double result = 0; for (int i = 5; i < 50; i = i + 5) { for (int j = i + 5; j < 50; j = j + 5) { optResult = OptimStart(i,j); if (optResult > result) { result = optResult; } } }

воскресенье, 1 марта 2020 г.

Немного вопросов о многопоточности

#java #многопоточность


Сам не первый год пишу на java, но лишь в рамках хобби, с многопоточностью приходится
не так часто работать. 

Заинтересовали несколько вопросов, а имённо:


Допустим, мне нужно писать и вычитывать примитивный тип на каком-то объекте (из разных
потоков). Допустим, поля открытые (public). 

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

Также справка java уверяет, что чтение и запись примитивных типов, исключая long
и double, тоже являются атомарными. Но как быть с этими двумя? Как я понимаю, их нужно
лишь пометить как volatile. Но есть ли какие-то побочные эффекты, или этот модификатор
только и делает, что гарантирует атомарность операций над переменной?

Хорошо. А теперь, к примеру, я хочу читать и писать поля через геттеры-сеттеры. Если
мои геттеры-сеттеры не меняют состояние каких-либо внешних данных, а лишь читают и
пишут переменную, достаточно ли мне пометить такую переменную как volatile, или же
нужно отмечать геттер/сеттер как syncronized? (При условии, что доступ к переменным
возможен лишь через геттер и сеттер).
Работа с объектами. А конкретно, с объектами, что не меняют состояния. 

К примеру, String. Насколько мне известно, данный тип никогда не меняет состояния,
а все его методы, возвращающие строку, возвращают новый экземпляр. 

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


Благодарю за терпение того, кто сможет на эти вопросы ответить)
    


Ответы

Ответ 1



Важно, что volatile гарантирует видимость. Изменения volatile-поля в одном потоке будут сразу видны в другом. Если вам нужно потокобезопасно выполнить более сложную процедуру, чем присвоение - воспользуйтесь synchronized или j.u.c.Atomic*-типами, чтобы избежать возникновения гонок (race condition). Пример: любой многопоточный счетчик // Плохо private volatile long counter; public long getId() { counter += 1; return counter; } // Хорошо private volatile AtomicLong counter = new AtomicLong(0); public long getId() { return counter.incrementAndGet(); } // Тоже хорошо private volatile long counter; public synchronized long getId() { counter += 1; return counter; } Помните, что synchronized - это всегда явная блокировка (но она довольно дёшева при низкой конкуренции), а Atomic-примитивы используют бесконечный цикл с CAS-инструкциями и могут хорошо оптимизироваться. Объекты с неизменяемым состоянием (например, все поля такого объекта final) назвают immutable объектами, и с ними безопасно работать из разных потоков. Но если внутреннее состояние объекта содержит изменяемый какой-то компонент и он как-то может быть доступен снаружи - придется дополнительно думать о потокобезопасности. Пример: // потокобезопасный класс public final class SafeFoo { private final Date date = new Date(); } // не потокобезопасный класс public final class UnsafeFoo { private final Date date = new Date(); public Date getDate() { return date; } } Первый класс полностью потокобезопасен. Второй - нет. Метод getDate() осуществляет публикацию внутреннего состояния внешнему коду (unsafe publication). И где-то снаружи можно сделать так, изменив внутренне состояние: unsafeFoo.getDate().setHours(0); Защититься от этого можно, например созданием копии: // потокобезопасный класс public final class SafeAgainFoo { private final Date date = new Date(); public Date getDate() { return date.clone(); } }