Сделайте fan-out RPC по поставщикам с дедлайном context и агрегацией
Реализуйте aggregateSearch, делающую fan-out по одному конкурентному SearchRPC(ctx, supplier, hotels) на поставщика и агрегирующую все результаты в map[hotel][]price. Требования: пробросьте ctx в каждый RPC, чтобы таймаут отменял вызовы в полёте, а не ждал их; собирайте результаты через канал (без общей мапы под mutex); закрывайте этот канал только после завершения всех горутин, никогда из отправителя.
func aggregateSearch(ctx context.Context, md []SearchMeta) map[string][]float64 {
// ваш код здесь
return nil
}
Допишите реализацию.
Запускают по одной goroutine на поставщика, каждая вызывает SearchRPC(ctx, …) и шлёт результаты в буферизованный канал; sync.WaitGroup и закрывающая goroutine делают close после завершения всех. Агрегируют, проходя range по каналу, и пробрасывают ctx, чтобы таймаут отменял все RPC в полёте, а не ждал их.
- ✗Создавать context через
context.WithTimeout, но не вызывать егоcancel— таймер течёт до дедлайна - ✗Закрывать канал результатов из goroutine-отправителя, а не после
wg.Wait, что даёт отправку в закрытый канал - ✗Не пробрасывать
ctxвSearchRPC, из-за чего дедлайн так и не отменяет медленного поставщика
- →Почему
closeканала должен идти послеwg.Wait, и почему в отдельной goroutine? - →Как пакет
errgroupсWithContextсократил бы этот код fan-out и агрегации?
Скелет для реализации
context.WithTimeout возвращает два значения (ctx, cancel); забытый cancel течёт таймером. Канал результатов закрывает отдельная goroutine после wg.Wait, а не отправитель. ctx пробрасывается в каждый SearchRPC, чтобы дедлайн отменял медленные вызовы.
ctx, cancel := context.WithTimeout(context.Background(), callTimeout)
defer cancel()
inverted := make(map[string][]string) // supplier -> []hotel
for _, m := range md {
inverted[m.Supplier] = append(inverted[m.Supplier], m.Hotel)
}
out := make(chan []SearchResult, len(inverted)) // буфер по числу отправителей
var wg sync.WaitGroup
for sup, hotels := range inverted {
wg.Add(1)
go func(sup string, hotels []string) {
defer wg.Done()
res, err := SearchRPC(ctx, sup, hotels) // ctx отменяет медленный RPC
if err != nil {
return
}
select {
case out <- res:
case <-ctx.Done(): // не блокируемся, если получатель ушёл
}
}(sup, hotels)
}
go func() { wg.Wait(); close(out) }() // close только после всех Done
prices := make(map[string][]float64)
for batch := range out { // агрегируем по мере прихода
for _, r := range batch {
prices[r.Hotel] = append(prices[r.Hotel], r.Price)
}
}