Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions concurrency/easy/basic-select.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,13 @@ import (

// Что выведет?

// Ответ:

// waited for 1 sec
// waited for 1 sec
// третий раз наверное не успеет выполниться, так как контекст закончится раньше, потому что опереации, которые не входят в select тоже займут какое-то небольшое время и третья "wait for 1 sec" не успеет и будет:
// deadline exceeded

func main() {
timeout := 3 * time.Second
ctx, cancel := context.WithTimeout(context.Background(), timeout)
Expand Down
56 changes: 55 additions & 1 deletion concurrency/easy/filtering-nums.go
Original file line number Diff line number Diff line change
@@ -1,4 +1,58 @@
// Есть горутина, которая генерирует случайные числа от 1 до N и отправляет в канал A.
// Вторая горутина читает из A и отправляет в канал B только чётные числа.
// Третья горутина читает из B и выводит числа.
// Можно с WaitGroup, можно с <-time.After()
// Можно с WaitGroup, можно с <-time.After()
package main

import (
"fmt"
"math/rand"
"sync"
)

var (
wg = &sync.WaitGroup{}
N = 100
)

func gorutine1(n int, A chan<- int) {
defer wg.Done()
defer close(A)

for i := 0; i < 10; i++ {
rnd := rand.Intn(n) + 1
A <- rnd
}

}

func gorutine2(A <-chan int, B chan<- int) {
defer wg.Done()
defer close(B)

for num := range A {
if num%2 == 0 {
B <- num
}
}
}

func gorutine3(B <-chan int) {
defer wg.Done()

for sorted_nums := range B {
fmt.Println(sorted_nums)
}
}

func main() {
A := make(chan int)
B := make(chan int)

wg.Add(3)
go gorutine1(N, A)
go gorutine2(A, B)
go gorutine3(B)

wg.Wait()
}
15 changes: 10 additions & 5 deletions concurrency/easy/generator-with-context.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,19 +2,21 @@
package main

import (
"context"
"fmt"
)

// Есть функция generate(), которая генерит числа
// Функция использует канал отмены. Переделать на контекст.
func generate(cancel <-chan struct{}, start int) <-chan int {
func generate(ctx context.Context, start int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for i := start; ; i++ {
select {
case out <- i:
case <-cancel:
// при срабатывании cancel() выполнится этот кейс и мы выйдем из функции
case <-ctx.Done():
return
}
}
Expand All @@ -23,13 +25,16 @@ func generate(cancel <-chan struct{}, start int) <-chan int {
}

func main() {
cancelCh := make(chan struct{})
// добавил контекст с отменой
ctx, cancel := context.WithCancel(context.Background())

generated := generate(cancelCh, 11)
// передаем теперь контекст, а не канал
generated := generate(ctx, 11)
for num := range generated {
fmt.Print(num, " ")
if num > 14 {
break
// как только наше условее "true" - вызываем cancel()
cancel()
}
}
fmt.Println()
Expand Down
32 changes: 32 additions & 0 deletions concurrency/easy/merge-channels.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ package main

import (
"fmt"
"sync"
"time"
)

Expand All @@ -23,6 +24,37 @@ func generateInRange(start, stop int) <-chan int {

func merge(channels ...<-chan int) <-chan int {
//TODO

// вейтгрупа для того, чтобы успеть считать все данные из каналов, а после закрыть нащ результирующий
wg := &sync.WaitGroup{}

// создаем и инициализируем наш результирующий канал
resCh := make(chan int)

// итерируемся по слайсу каналов (не уверен, что по слайсу, но поидее так и должно быть)
for _, ch := range channels {
// инкриментим счетчик
wg.Add(1)
go func(ch <-chan int) {
// читаем из канала
for item := range ch {
// пишем в наш результирующий канал канал
resCh <- item
}
// декриментим счетчик
wg.Done()
}(ch)
}

// небольшая махинация для того, чтоб:
// 1. не закрыть канал слишком рано, а потом горутиной писать в закрытый канал -> panic(),
// 2. вообще закрыть в целом, чтобы не словить deadlock, потому что "for val := range merged" будет ждать, а мы уже не пишем -> дедлок
go func() {
wg.Wait()
close(resCh)
}()

return resCh
}

func main() {
Expand Down
31 changes: 30 additions & 1 deletion concurrency/easy/time-after.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,15 +2,24 @@ package main

import (
"errors"
"fmt"
"time"
)

// Написать функцию after() — аналог time.After():

// возвращает канал, в котором появится значение
// через промежуток времени dur

func after(dur time.Duration) <-chan time.Time {
// ..
ch := make(chan time.Time, 1)
go func() {
defer close(ch)

time.Sleep(dur)
ch <- time.Time{}
}()
return ch
}

func withTimeout(fn func() int, timeout time.Duration) (int, error) {
Expand All @@ -29,3 +38,23 @@ func withTimeout(fn func() int, timeout time.Duration) (int, error) {
return 0, errors.New("timeout")
}
}

func main() {
// 1 <nil> (успевает выполниться работа)
fmt.Println(withTimeout(
func() int {
time.Sleep(time.Second * 1) // работает 1 сек
return 1
},
time.Second*1+time.Millisecond*200, // таймаут 1.2 сек
))

// 0 timeout (не успевает выполниться работа)
fmt.Println(withTimeout(
func() int {
time.Sleep(time.Second*1 + time.Millisecond*100) // работает 1.1 сек
return 1
},
time.Second*1, // таймаут 1 сек
))
}
23 changes: 22 additions & 1 deletion concurrency/easy/unpredictable-func.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
package main

import (
"context"
"fmt"
"math/rand"
"time"
Expand All @@ -27,7 +28,27 @@ func unpredictableFunc() int64 {
// Сигнатуру функцию обёртки менять можно.

func predictableFunc() int64 {
return unpredictableFunc()
start := time.Now()
var result int64
done := make(chan struct{})

ctx, _ := context.WithTimeout(context.Background(), time.Second)

go func(int64) {
result = unpredictableFunc()
close(done)
}(result)

select {
case <-done:
since := time.Since(start)
fmt.Println("Work my func:", since, "and my result:", result)
return result
case <-ctx.Done():
since := time.Since(start)
fmt.Println("Work my func:", since)
return 0
}
}

func main() {
Expand Down
6 changes: 0 additions & 6 deletions concurrency/hard/async-queue.go

This file was deleted.

28 changes: 28 additions & 0 deletions concurrency/hard/get-hotels-concurrent.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
package main

import (
"fmt"
"sync"
"time"
)

Expand All @@ -17,6 +19,32 @@ func main() {
hotelIDs := getHotels()

// Код здесь. Остальные функции и их сигнатуры менять нельзя. Допускается использовать функции-обертки.

// Создаю результирующий в который мы будем писать и из него же читать
resultCh := make(chan SearchResult)
// Чтоб закрыть канал после завершения работы горутин
wg := &sync.WaitGroup{}

// Читаем канал
for id := range hotelIDs {
wg.Add(1)

// Ищем инфу и пишем в канал
go func(id int) {
resultCh <- search(id)
wg.Done()
}(id)
}
// Закрываем канал
go func() {
wg.Wait()
close(resultCh)
}()

// Читаем результирующий канал
for i := range resultCh {
fmt.Println(i)
}
}

func search(hotelID int) SearchResult {
Expand Down
86 changes: 86 additions & 0 deletions concurrency/hard/task1/async-queue.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
// У нас есть поток задач которые нужно обрабатывать.
// Напишите библиотеку, которая будет принимать на вход задачи и асинхронно обрабатывать их.
// Работать обработчик должен в оперативной памяти.
// В один момент времени выполняется не более N задач,
// и не более М задач могут быть поставлены в очередь.
// Если в очереди нет места, возвращаем ошибку.

package main

import (
"fmt"
"math/rand"
"time"

taskqueue "github.com/f0xg0sasha/gotodev/tree/main/concurrency/hard/task1/asyncqueue"
)

func main() {
// Создаем очередь: 2 параллельных задачи, буфер на 6 задачи
q := taskqueue.New(3, 6)
q.Run() // Запускаем обработчики

// Функция задачи
createTask := func(id int, delay time.Duration) taskqueue.Task {
return func() {
fmt.Printf("Task %d started\n", id)
time.Sleep(delay)
fmt.Printf("Task %d completed: %d \n", id, id*id)
}
}

// Добавляем задачи
for i := 1; i <= 15; i++ {
err := q.AddTask(createTask(i, time.Duration(rand.Intn(3))*500*time.Millisecond))
if err != nil {
fmt.Printf("Error adding task %d: %v\n", i, err)
} else {
fmt.Printf("Task %d added to queue\n", i)
}
}

// Ждем завершения всех задач
q.Wait()
fmt.Println("All tasks completed")
}

// $ go run async-queue.go

// Task 1 added to queue
// Task 2 added to queue
// Task 3 added to queue
// Task 4 added to queue
// Task 5 added to queue
// Task 6 added to queue
// Task 7 added to queue
// Task 8 added to queue
// Task 9 added to queue
// Task 1 started
// Task 2 started
// Task 2 completed: 4
// Task 4 started
// Task 4 completed: 16
// Task 5 started
// Task 3 started
// Error adding task 10: queue is full
// Task 11 added to queue
// Task 12 added to queue
// Error adding task 13: queue is full
// Error adding task 14: queue is full
// Error adding task 15: queue is full
// Task 5 completed: 25
// Task 6 started
// Task 6 completed: 36
// Task 7 started
// Task 1 completed: 1
// Task 8 started
// Task 3 completed: 9
// Task 9 started
// Task 9 completed: 81
// Task 11 started
// Task 7 completed: 49
// Task 12 started
// Task 8 completed: 64
// Task 11 completed: 121
// Task 12 completed: 144
// All tasks completed
Loading