Roman Kryvolapov Engineering Blog

Java / Kotlin Core — Багатопотоковість

Які існують основні оператори в RxJava

У RxJava (Reactive Extensions for Java) існує безліч операторів, які дозволяють працювати з потоками даних (Observable) асинхронно та реактивно. Ці оператори допомагають створювати, трансформувати, фільтрувати, комбінувати та обробляти потоки даних.

Ось основні категорії операторів і приклади їх використання:

Оператори створення (Creating Operators):
Ці оператори використовуються для створення потоків (Observable).

just():
Створює потік з одним або кількома елементами.

Observable.just("Hello", "RxJava")
.subscribe(System.out::println);

fromIterable():
Перетворює колекцію на потік даних.

List<String> items = Arrays.asList("One", "Two", "Three");
Observable.fromIterable(items)
.subscribe(System.out::println);

create():
Створює потік вручну, з використанням логіки емісії.

Observable.create(emitter -> {
emitter.onNext("Data");
emitter.onComplete();
}).subscribe(System.out::println);

range():
Емітує послідовність чисел у вказаному діапазоні.

Observable.range(1, 5)
.subscribe(System.out::println);

Оператори перетворення (Transforming Operators):
Ці оператори перетворюють дані в потоці.

map():
Перетворює кожен елемент потоку.

Observable.just(1, 2, 3)
.map(item -> item * 2)
.subscribe(System.out::println);

flatMap():
Перетворює кожен елемент вихідного потоку на новий потік і об'єднує їх в один загальний потік (асинхронно).

Observable.just(1, 2, 3)
.flatMap(item -> Observable.just(item * 2))
.subscribe(System.out::println);

concatMap():
Схожий на flatMap(), але виконує перетворення елементів послідовно (синхронно), у порядку їх надходження.

Observable.just(1, 2, 3)
.concatMap(item -> Observable.just(item * 10))
.subscribe(System.out::println); // 10, 20, 30

switchMap():
Перетворює елемент на новий потік, але якщо надходить новий елемент, попередній потік припиняє емісію, і замість нього емітується новий потік. Це корисно, коли потрібно завжди працювати лише з останнім потоком.

Observable.just(1, 2, 3)
.switchMap(item -> Observable.just(item * 10))
.subscribe(System.out::println); // Лише останній елемент з нового потоку

groupBy():
Розділяє потік на кілька потоків (груп) на основі певного критерію.

Observable.just(1, 2, 3, 4, 5, 6)
.groupBy(item -> item % 2 == 0 ? "Even" : "Odd")
.subscribe(groupedObservable -> {
groupedObservable.subscribe(item ->
System.out.println(groupedObservable.getKey() + ": " + item)
);
}); // Odd: 1 Even: 2 Odd: 3 Even: 4 Odd: 5 Even: 6

flatMapIterable():
Перетворює кожен елемент на Iterable та емітує кожен елемент цього Iterable.

Observable.just(1, 2, 3)
.flatMapIterable(item -> Arrays.asList(item * 10, item * 100))
.subscribe(System.out::println); // Результат: 10, 100, 20, 200, 30, 300

cast():
Перетворює елементи в потоках до певного типу.

Observable<Object> observable = Observable.just(1, "Hello", 3.14);
observable.cast(String.class)
.subscribe(System.out::println, Throwable::printStackTrace); // Помилка приведення типів

buffer():
Збирає елементи в буфер (наприклад, по кілька штук) та емітує їх як списки.

Observable.range(1, 10)
.buffer(3)
.subscribe(System.out::println); // Виводить списки [1, 2, 3], [4, 5, 6], ...

scan():
Застосовує акумуляторну функцію до кожного елемента та повертає кожен проміжний результат.

Observable.just(1, 2, 3)
.scan((acc, item) -> acc + item)
.subscribe(System.out::println); // 1, 3, 6

reduce():
Схожий на scan(), але повертає лише кінцевий результат застосування акумуляторної функції.

Observable.just(1, 2, 3)
.reduce((accumulator, item) -> accumulator + item)
.subscribe(System.out::println); // Результат: 6

window():
Розбиває потік на окремі вікна (підпотоки) з певною кількістю елементів або часом.

Observable.range(1, 10)
.window(3)
.subscribe(window -> {
System.out.println("New Window:");
window.subscribe(System.out::println);
});
New Window:
1
2
3
New Window:
4
5
6
New Window:
7
8
9
New Window:
10

toList():
Перетворює потік даних на один список.

Observable.just(1, 2, 3)
.toList()
.subscribe(System.out::println); // Результат: [1, 2, 3]

toMap():
Перетворює потік даних на Map, де ключі генеруються за допомогою функції.

Observable.just("apple", "banana", "cherry")
.toMap(fruit -> fruit.charAt(0))
.subscribe(System.out::println); // Результат: {a=apple, b=banana, c=cherry}

startWith():
Додає елементи на початок потоку перед емісією оригінальних елементів.

Observable.just(2, 3, 4)
.startWith(1)
.subscribe(System.out::println); // Результат: 1, 2, 3, 4

repeat():
Повторює потік вказану кількість разів.

Observable.just(1, 2, 3)
.repeat(2)
.subscribe(System.out::println); // Результат: 1, 2, 3, 1, 2, 3

delay():
Затримує емісію елементів на певний час.

Observable.just(1, 2, 3)
.delay(2, TimeUnit.SECONDS)
.subscribe(System.out::println); // Емітує елементи через 2 секунди

debounce():
Емітує елемент лише якщо минув певний час з останньої емісії (популярно для фільтрації подій, що швидко надходять, наприклад, при роботі з користувацькими введеннями).

Observable.create(emitter -> {
emitter.onNext(1);
Thread.sleep(300);
emitter.onNext(2);
Thread.sleep(100);
emitter.onNext(3);
Thread.sleep(400);
emitter.onComplete();
})
.debounce(200, TimeUnit.MILLISECONDS)
.subscribe(System.out::println); // Емітує лише 1 і 3 (ігнорує 2, оскільки надто швидко після 1)

throttleFirst():
Емітує перший елемент з потоку за заданий інтервал часу та ігнорує решту елементів у цей період.

Observable.interval(100, TimeUnit.MILLISECONDS)
.throttleFirst(1, TimeUnit.SECONDS)
.subscribe(System.out::println);

throttleLast() (або sample()):
Емітує останній елемент, який був згенерований у заданий інтервал часу.

Observable.interval(100, TimeUnit.MILLISECONDS)
.throttleLast(1, TimeUnit.SECONDS)
.subscribe(System.out::println);

Оператори фільтрації (Filtering Operators):
Ці оператори використовуються для фільтрації даних у потоці.

filter():
Пропускає лише ті елементи, які задовольняють умову.

Observable.just(1, 2, 3, 4, 5)
.filter(item -> item % 2 == 0)
.subscribe(System.out::println); // 2, 4

distinct():
Пропускає лише унікальні елементи.

Observable.just(1, 2, 2, 3, 3, 3)
.distinct()
.subscribe(System.out::println); // 1, 2, 3

take():
Емітує лише вказану кількість елементів.

Observable.just(1, 2, 3, 4)
.take(2)
.subscribe(System.out::println); // 1, 2

skip():
Пропускає вказану кількість елементів, а потім емітує решту.

Observable.just(1, 2, 3, 4)
.skip(2)
.subscribe(System.out::println); // 3, 4

debounce():
Емітує елемент лише якщо минув певний час з останньої емісії.

Observable.just(1, 2, 3, 4)
.debounce(100, TimeUnit.MILLISECONDS)
.subscribe(System.out::println);

*Оператори комбінування (Combining Operators):* Ці оператори дозволяють об'єднувати кілька потоків даних.

merge():
Об'єднує кілька потоків в один, чергуючи елементи.

Observable.merge(
Observable.just("A", "B"),
Observable.just("1", "2")
).subscribe(System.out::println); // A, 1, B, 2

zip():
Комбінує елементи з кількох потоків, використовуючи функцію.

Observable.zip(
Observable.just("A", "B"),
Observable.just("1", "2"),
(letter, number) -> letter + number
).subscribe(System.out::println); // A1, B2

concat():
Емітує елементи з одного потоку, потім переходить до наступного.

Observable.concat(
Observable.just("A", "B"),
Observable.just("1", "2")
).subscribe(System.out::println); // A, B, 1, 2

combineLatest():
Емітує останній елемент з кожного потоку щоразу, коли один з потоків емітує новий елемент.

Observable.combineLatest(
Observable.just("A", "B"),
Observable.just("1", "2", "3"),
(letter, number) -> letter + number
).subscribe(System.out::println); // B1, B2, B3

Оператори для обробки помилок (Error Handling Operators):
Ці оператори допомагають керувати помилками в потоці.

onErrorReturn():
Повертає елемент, якщо сталася помилка.

Observable.just(1, 2, 0)
.map(item -> 10 / item)
.onErrorReturn(e -> -1)
.subscribe(System.out::println); // 10, 5, -1

retry():
Повторює виконання потоку при помилці.

Observable.just(1, 2, 0)
.map(item -> 10 / item)
.retry(2)
.subscribe(System.out::println, Throwable::printStackTrace);

onErrorResumeNext():
Продовжує емітувати елементи з іншого потоку у випадку помилки.

Observable.just(1, 2, 0)
.map(item -> 10 / item)
.onErrorResumeNext(Observable.just(-1))
.subscribe(System.out::println);

Оператори роботи з часом (Time-based Operators):
Ці оператори дозволяють керувати часом у потоках.

interval():
Емітує елементи через регулярні проміжки часу.

Observable.interval(1, TimeUnit.SECONDS)
.subscribe(System.out::println);

timer():
Емітує один елемент після вказаної затримки.

Observable.timer(2, TimeUnit.SECONDS)
.subscribe(System.out::println);

Оператори підключення (Connectable Operators):
Вони дозволяють робити холодні потоки гарячими, тобто починати емісію даних лише після виклику спеціального методу.

publish():
Перетворює потік на "гарячий", але емісія даних починається лише після виклику connect().

ConnectableObservable<Long> connectable = Observable.interval(1, TimeUnit.SECONDS).publish();
connectable.connect(); // Запускає емісію даних

Які бувають і як працюють синхронізовані колекції

CopyOnWrite колекції:
усі операції зі зміни колекції (add, set, remove) призводять до створення нової копії внутрішнього масиву. Тим самим гарантується, що при проході ітератором по колекції не буде кинуто ConcurrentModificationException.
CopyOnWriteArrayList, CopyOnWriteArraySet

Scalable Maps:
покращені реалізації HashMap, TreeMap з кращою підтримкою багатопотоковості та масштабованості.
ConcurrentMap, ConcurrentHashMap, ConcurrentNavigableMap, ConcurrentSkipListMap, ConcurrentSkipListSet

Non-Blocking Queues:
потокобезпечні та неблокуючі імплементації Queue на зв'язаних нодах (linked nodes).
ConcurrentLinkedQueue, ConcurrentLinkedDeque

Blocking Queues:
BlockingQueue, ArrayBlockingQueue, DelayQueue, LinkedBlockingQueue, PriorityBlockingQueue, SynchronousQueue, BlockingDeque, LinkedBlockingDeque, TransferQueue, LinkedTransferQueue

Що означає volatile

Ключове слово volatile у Kotlin (і в Java) використовується для оголошення змінних, значення яких можуть бути змінені різними потоками. Воно гарантує, що при зміні змінної значення буде одразу ж записано в пам'ять і зчитано з пам'яті, а не буде закешоване в регістрах процесора, що може призвести до проблем синхронізації та видимості змін.
Ключове слово volatile може бути застосоване до змінних типу Boolean, Byte, Char, Short, Int, Long, Float, Double, а також до посилальних типів.

@Volatile
private var running: Boolean = false

У цьому прикладі змінна running буде оновлюватися в пам'яті одразу ж після зміни будь-яким потоком, і потоки, які використовують її значення, бачитимуть найостанніше значення в пам'яті.

Що означає synchronized

ключове слово в Kotlin (і в Java), яке використовується для синхронізації доступу до спільних ресурсів у багатопотокових застосунках. Коли кілька потоків намагаються отримати доступ до спільного ресурсу одночасно, можуть виникнути проблеми, наприклад, неоднозначність стану ресурсу або його пошкодження. synchronized дозволяє уникнути цих проблем, шляхом гарантованого одночасного доступу лише одному потоку до спільного ресурсу.

synchronized(lock) {
// блок коду, в якому виконується доступ до спільного ресурсу
}
@Synchronized
fun getCounter(): Int {
return counter
}

Корутини (coroutines) не можуть бути synchronized, оскільки synchronized є ключовим словом у Java, яке використовується для синхронізації доступу до спільних ресурсів між кількома потоками. Замість цього, для синхронізації доступу до спільних ресурсів між корутинами в Kotlin, слід використовувати інші механізми синхронізації, такі як атомарні змінні (atomic variables), блокування (locks) або м'ютекси (mutexes). Наприклад, можна використати м'ютекс зі стандартної бібліотеки Kotlin:

import kotlinx.coroutines.*
import kotlinx.coroutines.sync.Mutex
val mutex = Mutex()
fun main() = runBlocking {
launch {
mutex.withLock {
// код, який потрібно виконати синхронно
}
}
}

У цьому прикладі використовується м'ютекс Mutex зі стандартної бібліотеки Kotlin, який дозволяє заблокувати доступ до спільного ресурсу всередині блоку withLock. Це гарантує, що код, який знаходиться всередині блоку withLock, виконуватиметься лише однією корутиною в будь-який момент часу.

Які проблеми можуть бути в багатопотоковості в Java

Race condition:
ситуація, за якої результат виконання програми залежить від того, які потоки виконуватимуться швидше або пізніше.

var count = 0
fun main() {
Thread {
for (i in 1..100000) {
count++
}
}.start()
Thread {
for (i in 1..100000) {
count++
}
}.start()
Thread.sleep(1000)
println("Count: $count")
}

У цьому прикладі створюються два потоки, кожен з яких збільшує змінну count на 1 однакову кількість разів. Після виконання цих потоків у головному потоці виводиться значення змінної count. Однак, оскільки два потоки можуть виконуватися паралельно, результат залежить від того, який потік закінчить роботу першим, і може бути несподіваним.

Наприклад, при одному запуску програми результат може бути 199836, а при іншому — 200000. Це відбувається через те, що змінна count використовується не атомарно (не захищена механізмами синхронізації), і два потоки можуть змінювати її значення одночасно. У результаті, значення змінної можуть перезаписуватися, і остаточне значення може бути меншим, ніж очікувалося.

Щоб уникнути Race condition у таких випадках, необхідно використовувати синхронізацію та механізми захисту, такі як блокування (lock) та атомарні змінні (atomic variables), які дозволяють гарантувати коректність виконання операцій та уникнути несподіваних результатів.

Deadlock:
виникає, коли два або більше потоки блокуються, чекаючи одне на одного, щоб звільнити ресурси, необхідні для продовження виконання.

Неправильне використання synchronized: synchronized може бути використаний неправильно, що може призвести до неправильного порядку виконання або заблокованих потоків.

Щоб уникнути цих проблем, у Kotlin можна використовувати засоби синхронізації, такі як mutex, lock, atomic variables та інші, а також користуватися сучасними практиками багатопотокового програмування, такими як використання immutable data structures та обмеження зміни спільних даних лише всередині критичних секцій.

class Resource(private val name: String) {
@Synchronized fun checkResource(other: Resource) {
println("$this: checking ${other.name}")
Thread.sleep(1000)
other.checkResource(this)
}
}
fun main() {
val resource1 = Resource("Resource 1")
val resource2 = Resource("Resource 2")
Thread {
resource1.checkResource(resource2)
}.start()
Thread {
resource2.checkResource(resource1)
}.start()
}

У цьому прикладі створюються два об'єкти класу Resource, кожен з яких синхронізований за допомогою анотації @Synchronized. Далі створюються два потоки, кожен з яких викликає метод checkResource() для різних ресурсів у різних порядках. Таким чином, якщо перший потік заблокує resource1, а другий потік заблокує resource2, то обидва потоки чекатимуть одне на одного і не зможуть завершити виконання, що призведе до Deadlock. Приклад Deadlock у Kotlin можна уникнути, якщо взаємодія між потоками відбувається з використанням спільних блокувань, наприклад, за допомогою synchronized або lock() зі стандартної бібліотеки Kotlin. Також можна використовувати методи wait() та notify() для керування потоками та уникнення блокувань.

Livelock:
ситуація, за якої два або більше потоки продовжують виконувати дії для уникнення блокування, але не можуть завершити свою роботу. У результаті, вони перебувають у нескінченному циклі, споживаючи все більше ресурсів і не виконуючи жодної корисної роботи

data class Person(val name: String, val isPolite: Boolean = true) {
fun greet(other: Person) {
while (true) {
if (isPolite) {
println("$name: After you, ${other.name}")
Thread.sleep(1000)
if (other.isPolite) {
break
}
} else {
println("$name: No, please, after you, ${other.name}")
Thread.sleep(1000)
if (!other.isPolite) {
break
}
}
}
}
}
fun main() {
val john = Person("John", true)
val jane = Person("Jane", false)
Thread {
john.greet(jane)
}.start()
Thread {
jane.greet(john)
}.start()
}

У цьому прикладі створюються два об'єкти класу Person, кожен з яких може бути ввічливим або неввічливим залежно від значення властивості isPolite. Потім створюються два потоки, кожен з яких викликає метод greet() для іншого об'єкта. Метод greet() використовує цикл while, щоб перевіряти, чи є інший об'єкт ввічливим, і, залежно від цього, продовжувати або зупинятися.

Якщо обидва об'єкти будуть ввічливими, то вони по черзі пропонуватимуть одне одному піти першим і не зможуть закінчити свою розмову. Якщо ж обидва об'єкти будуть неввічливими, то вони відмовлятимуться йти першими і також не зможуть завершити розмову. Таким чином, обидва потоки продовжуватимуть свою роботу в нескінченному циклі, не виконуючи жодної корисної роботи, що є прикладом Livelock.

Прикладу Livelock можна уникнути, наприклад, за допомогою використання таймерів та обмеження часу виконання операцій, щоб потоки могли завершити свою роботу та уникнути нескінченного циклу. Також можна використовувати синхронізацію та блокування, щоб забезпечити коректну взаємодію між потоками.

Що таке Lock / ReentrantLock

Lock:
це механізм синхронізації, який використовується для запобігання конкуренції за доступ до спільних ресурсів між кількома потоками в багатопотоковому середовищі. Коли потік використовує lock, він отримує ексклюзивний доступ до ресурсу, пов'язаного із замком, і інші потоки, які намагаються отримати доступ до цього ресурсу, блокуються, доки перший потік не звільнить lock.

ReentrantLock:
це реалізація інтерфейсу Lock у Java, яка дозволяє потоку захоплювати та звільняти lock багаторазово. Він називається “перевикористовуваним”, тому що він дозволяє потоку захоплювати lock знову, якщо цей потік уже має доступ до заблокованого ресурсу.
ReentrantLock дозволяє контролювати механізм блокування ретельніше, ніж за допомогою синхронізованих блоків, оскільки він надає деякі додаткові функції, такі як можливість призупиняти та відновлювати потоки, які очікують доступу до заблокованого ресурсу, можливість встановлення тайм-ауту на очікування доступу до ресурсу та можливість використання кількох умовних змінних для керування доступом до ресурсу.

Які є стратегії Backpressure

Buffer:
буферизує всі елементи, які були відправлені, доки споживач не буде готовий їх обробити. Ця стратегія може призвести до нестачі пам'яті, якщо виробник генерує елементи занадто швидко або споживач обробляє елементи занадто повільно.

Drop:
відкидає елементи, які не можуть бути оброблені споживачем. Ця стратегія не гарантує, що всі елементи будуть оброблені, але дозволяє уникнути нестачі пам'яті.

Latest:
зберігає лише останній елемент і відкидає всі інші. Ця стратегія підходить для сценаріїв, коли лише останній елемент має значення.

Error:
сигналізує про помилку, якщо потік не може обробити дані.

Missing:
якщо потік не може обробляти дані, то він буде просто пропускати їх, без попередження.

Чим відрізняються hot і cold Observables

Cold Observable:
Не розсилає об'єкти, доки на нього не підписався хоча б один підписник;
Якщо observable має кількох підписників, то він розсилатиме всю послідовність об'єктів кожному підписнику.

Hot Observable:
Розсилає об'єкти, коли вони з'являються, незалежно від того, чи є підписники;
Кожен новий підписник отримує лише нові об'єкти, а не всю послідовність.

Чим volatile відрізняється від atomic

Ключове слово volatile і класи Atomic* у Java використовуються для забезпечення безпеки потоків при доступі до спільної пам'яті.

volatile:
гарантує, що значення змінної завжди буде прочитано зі спільної пам'яті, а не з кешу потоку, що запобігає помилкам синхронізації. Крім того, запис у volatile змінну теж записується безпосередньо у спільну пам'ять, а не в кеш потоку.

Класи Atomic:
забезпечують атомарність операцій зі змінними. Тобто, вони гарантують, що операції читання та запису будуть виконані як єдина, неподільна дія. Класи Atomic реалізовані з використанням механізмів апаратної підтримки, тому вони можуть бути ефективнішими в деяких випадках, ніж використання volatile.

Загалом, якщо потрібно лише забезпечити безпеку потоків при доступі до спільної пам'яті, то можна використовувати volatile. Якщо ж необхідно виконати атомарну операцію на змінній, то потрібно використовувати класи Atomic*.

Що таке відкладені (Deferred) корутини в Kotlin

Один з механізмів роботи з асинхронними операціями в Kotlin з використанням корутин. Вони являють собою об'єкти, які представляють результат виконання асинхронної операції.

Відкладені корутини створюються за допомогою функції async і можуть бути використані для виконання довгих операцій у фоновому потоці, при цьому основний потік застосунку залишається вільним.
Щойно відкладена корутина завершує свою роботу, її результат може бути отриманий за допомогою функції await(). На відміну від функції join(), яка блокує потік, що викликає, до завершення виконання корутини, await() не блокує потік, що викликає, а повертає значення лише тоді, коли воно готове.
Ось приклад використання відкладених корутин у Kotlin:

fun loadDataAsync(): Deferred<List<Data>> = GlobalScope.async {
// завантаження даних з мережі або бази даних у фоновому потоці
}
fun displayData() {
GlobalScope.launch {
// запускаємо корутину для завантаження даних
val deferredData = loadDataAsync()
// виконаємо деякі дії в основному потоці
// отримуємо результат виконання відкладеної корутини
val data = deferredData.await()
// обробляємо дані
}
}

Що таке ліниві (LAZY) корутини в Kotlin

LAZY корутини це спосіб створення відкладених корутин у Kotlin за допомогою функції lazy та async з бібліотеки kotlinx.coroutines.
Основна відмінність відкладених корутин, створених за допомогою функції async, від LAZY корутин полягає в моменті створення об'єкта корутини. Відкладені корутини створюються негайно, коли викликається функція async, тоді як LAZY корутини створюються лише в той момент, коли до них звертаються вперше.
Як правило, LAZY корутини використовуються у випадках, коли немає необхідності одразу починати виконання корутини, наприклад, коли ми не знаємо, чи буде її виконання взагалі необхідним. Це дозволяє уникнути зайвих ресурсовитрат і пришвидшити роботу застосунку.
Ось приклад створення LAZY корутини:

val lazyCoroutine: Lazy<Deferred<String>> = lazy {
GlobalScope.async {
// виконання корутини у фоновому потоці
"Hello from coroutine!"
}
}
fun main() {
// виконання дій в основному потоці
// отримання результату виконання корутини
val result = runBlocking { lazyCoroutine.value.await() }
// обробка результату
println(result)
}
// або
fun main() = runBlocking<Unit> {
val lazyCoroutine = launch(start = CoroutineStart.LAZY) {
println("Coroutine is executing")
}
// виконання інших дій в основному потоці
lazyCoroutine.start() // запуск корутини
// виконання інших дій в основному потоці
lazyCoroutine.join() // очікування завершення корутини
}

У цьому прикладі ми створюємо відкладену корутину за допомогою функції launch і передаємо їй параметр start = CoroutineStart.LAZY. Потім ми виконуємо якісь інші дії в основному потоці, і лише після цього ми викликаємо метод start() нашої відкладеної корутини, щоб запустити її виконання.
Важливо розуміти, що якщо ми не викликаємо метод start() нашої відкладеної корутини, то вона ніколи не буде виконана.
Ми також можемо використовувати параметр CoroutineStart.LAZY з функціями async та withContext. У цьому випадку створюється відкладена корутина, яка повертає результат обчислень. Ми можемо викликати await() на цій корутині, щоб отримати результат виконання.

Що таке Kotlin Channels

Компонент, наданий Kotlin Coroutines, який дозволяє асинхронно передавати значення між корутинами.
Канали (Channels) є більш високорівневою абстракцією, порівняно з примітивними синхронізованими об'єктами, такими як блокування та семафори. Вони дозволяють передавати значення між корутинами, підтримуючи при цьому коректний порядок і безпеку в багатопотоковому середовищі.

Kotlin Channels схожі на черги повідомлень і мають такі основні операції: відправлення (send) та отримання (receive). Також існує можливість закрити канал (close), що призводить до зупинки передачі даних.
Крім того, Kotlin Channels надають різні операції буферизації, такі як обмеження розміру буфера та часовий таймаут для отримання повідомлення.

Загалом, Kotlin Channels дозволяють ефективніше та безпечніше організувати асинхронну передачу даних у застосунку.

fun main() = runBlocking {
val channel = Channel<Int>()
launch {
for (x in 1..5) channel.send(x * x)
}
repeat(5) { println(channel.receive()) }
println("Done!")
}
// Очікуваний вивід: 1 4 9 16 25 Done!

цьому прикладі ми створили Channel за допомогою Channel(), який дозволяє передавати цілі числа. Потім ми запустили корутину, яка відправляє п'ять повідомлень у канал методом send(). У головному потоці ми отримуємо повідомлення з каналу методом receive() і виводимо їх на консоль.
Тут ми використали функцію runBlocking для запуску основного потоку. Це необхідно, оскільки функція launch запускає корутину в новому потоці

Що таке Mutex, Monitor, Semaphore

Semaphore:
тип блокування, обмежує кількість потоків, які можуть увійти в задану ділянку коду.

Mutex:
об'єкт для синхронізації потоків, прикріплений до кожного об'єкта в Java.
Може мати 2 стани — вільний і зайнятий. Станом м'ютекса не можна керувати напряму.
Якщо іншому потоку буде потрібен доступ до змінної, захищеної м'ютексом, то цей потік блокується доти, доки м'ютекс не буде звільнений.
Відрізняється від семафора тим, що лише потік, який ним володіє, може його звільнити.
У блоці коду, який позначений словом synchronized, відбувається захоплення м'ютекса.
Завданням м'ютекса є захист об'єкта від доступу до нього інших потоків, відмінних від того, який заволодів м'ютексом.
М'ютекс захищає дані від пошкодження в результаті асинхронних змін (стан гонки), однак при неправильному використанні можуть породжуватися інші проблеми, наприклад, взаємне блокування або подвійне захоплення.

Monitor:
високорівневий механізм взаємодії та синхронізації процесів, що забезпечує доступ до нероздільних ресурсів.
Створює захисний механізм для реалізації synchronized блоків.

Copyright: Roman Kryvolapov