Блокировка горутин (не deadlock)
Как заблокировать множество горутин, кроме одной? А затем быстро разблокировать их?
Суть задачи в том, что некоторое множество горутин могут одномоментно получить от API некоего сервиса одну и ту же ошибку.
И эту ошибку должна разрулить одна единственная горутина (взять из deque новый ключ и затем сделать запрос API, чтобы проверить его. Если ключ рабочий - разблокировать прочие горутины. Если нет - взять еще один ключ, или передать такую возможность другой горутине, а самой заблокироваться в ожидании). И так до получения от API кода 200.
Не получается придумать правильный алгоритм (подход) для такой логики, чтобы только одна горутина могла бы обращаться к очереди с ключами до тех пор, пока там вообще что-то есть или пока не найден рабочий ключ.
На текущий момент горутины после ошибки просто блокируются через sleep c последовательным увеличением (инкремент залочен мьютексом) переменной шага таймаута. Таким образом шанс разрулить получает та горутина, чей таймаут ожидания меньше прочих. Это работает (не очень хорошо, там еще приходится использовать цикл с условной переменной,
for trs.Suspended {
log.Debug("Suspended")
time.Sleep(1 * time.Second)
}
пока она true - горутина лочится в нем). Но такой подход мне совсем не нравится.
Ответы (2 шт):
- Вам в таком случае не нужны горутины (если всегда работает ТОЛЬКО ОДИН обработчик). Достаточно создать словарь вызовов соответствующего обработчика в однотипным объектах (лучше чтобы элементы словаря были интерфейсами -- можно отдать ссылки на любой тип, который удовлетворяет этому интерфейсу). Также замечу, что когда ОЧЕНЬ МНОГО горутин -- рантайм тратит существенное время на обработку состояния каналов (на моём рабочем ноуте 1200 каналов требует 5% процессорного времени на ровном месте -- пришлось сокращать эту басню).
mapRunner:=make(map[name]*IRunner, n)
sigName:=<-chanSig
mapRanner[sigName].Run()
Эта схема работы называется "селектор". 2) Если всё же вам РЕАЛЬНО НУЖНА АСИНХРОННОСТЬ -- сделайте словарь каналов с выходом в каждой горутине. Словарём каналов должна владеть управляющая горутина. Как только приходит сигнал с нужным параметром -- управляющая горутина засылает в соответствующий канал сигнал. В качестве словаря также можно использовать интерфейсы, но тогда указанные интерфейсы должны иметь методы для запихивания сигналов в свои каналы и вызовы закрытия этих каналов для прерывания работы горутин (это более универсально, но требует больше работы).
const(
SIG = 0
)
mapRunner:=make(map[name]*TRunner, n)
sigName:=<-chanSig
mapRunner[sigName]<-SIG
Эта схема работы называется "демультиплексор".
В итоге решение получилось такое. В main функции заранее создается канал на один тикет.
trs.Ticket = make(chan bool,1)
trs.TicketDelay = time.Duration(5000)
trs.TimeBlocking = time.Duration(1000)
trs.Ticket <- true
В который сразу пишется одно значение - билет-проходка. В горутинах, которые что-то запрашивают и могут получить ошибки добавлен select для чтения из канал Ticket.
Логика такая: получили ошибку, идем в if, если ошибка подходит, идем в select, кто-то первый успевает прочитать (достать) ticket, выходит из select'а, идет в switch - решает (пытается решить) проблему. Отдает тикет обратно, чтобы какая-нибудь другая горутина смогла им воспользоваться, если проблема не решилась (горутины опять получили ошибку).
Горутины, которым тикет не достался блокируются на канале-таймере (1 сек), затем возвращаются в цикл запросов и пробуют снова сделать запрос. Если опять ошибка, опять кто-то берет тикет, кто-то ждет, ну и так до победного....
...
defer done.Wg.Done()
TASKLOOP:
for {
select {
case <-trs.ctx.Done():
YAPI = &YandexAPIError{500, "The task was canceled", trs.ctx.Err()}
break TASKLOOP
default:
resp, err := trs.API.Translate(task.Lang, task.Text)
... какой-то код
// смотрим что за ошибка, если из списка - пробуем решить
if contains([]int{401, 402, 404, 405, 406},
resp.Code) {
select {
// канал на один тикет - кто успел взять, тот и решает проблему
case <-trs.Ticket:
break // выходим из select и идем в switch
// здесь блокируемся, если билетика нет, а потом идем делать новый запрос
case <-time.After(trs.TimeBlocking * time.Millisecond):
goto TASKLOOP
}
}
switch resp.Code {
.... много других кейсов
case 405, 406: // 'Session is invalid', 'Session has expired',
trs.mutex.Lock()
log.Errorf("[SESSION][task:%d][code:%d] EXPIRED or INVALID %s",
task.Id, resp.Code, trs.API.Id)
change := trs.ChangeSession(resp) // меняем session id
trs.mutex.Unlock()
// если получилось взять из очереди новый ключ\session_id
if change {
log.Warnf("[CHANGE SESSION] %s", trs.API.Id)
// немножко тормозим отдачу билетика в канал
time.Sleep(trs.TicketDelay * time.Millisecond)
trs.Ticket <- true // возвращаем билетик обратно
goto TASKLOOP // идем на новый виток запросов
} else {
// ну если ничего не вышло - в очереди больше нет ключей\session_id
YAPI = &YandexAPIError{resp.Code, resp.Message, err}
log.Errorf("[CHANGE SESSION] %s", NotFoundAPISessionID)
trs.ctx.Cancel() // отменяем все задания
break TASKLOOP // завершаем горутину
}
}
} // конец цикла
if YAPI != nil {
task.Err = YAPI
done.TaskChan <- task
log.Errorf("[TRANSLATE][%d] FAIL %#v", task.Id, YAPI)
} else {
log.Debugf("[TRANSLATE][%d] DONE!", task.Id)
}
Не знаю, насколько это решение простое и гениальное, но я до него долго доходил :-) Можно покритиковать. Всем спасибо за ответы.