Yandex / Яндекс (3й этап)
Актуальность: 4 кв 2025
Задача: Нам нужно передать данные из некоторого источника некоторому потребителю. При этом источник отдает данные небольшими пачками (~ десятки записей), а потребитель оптимальнее работает с крупными батчами (~тысячи записей). Реальный пример - поставка данных из очередей типа Kafka в базу Clickhouse. go/concurrency
Заголовок раздела «Задача: Нам нужно передать данные из некоторого источника некоторому потребителю. При этом источник отдает данные небольшими пачками (~ десятки записей), а потребитель оптимальнее работает с крупными батчами (~тысячи записей). Реальный пример - поставка данных из очередей типа Kafka в базу Clickhouse. go/concurrency»Источник
Заголовок раздела «Источник»- Условно бесконечный.
go/concurrency - Источник никогда не возвращает более MaxItems записей за один вызов Next.
go/concurrency - В рамках одной «сессии»(одного вызова функции Pipe) источник каждый раз возвращает новые данные на каждый вызов Next.
go/concurrency - Однако, после перезапуска источник начнет с прошлой «подтвержденной» позиции, задаваемой cookie. Поэтому каждое значение cookie, которое вернул вызов Next, после сохранения данных в приемнике, должно быть фиксировано вызовом Commit, причем строго в той же последовательности, в которой их вернул Next
network/http
Приемник
Заголовок раздела «Приемник»- Не может обработать более MaxItems за один раз.
go/concurrency
Задача: (код без описания) go/basics
Заголовок раздела «Задача: (код без описания) go/basics»const MaxItems = 9999
type Producer interface { // Next returns: // - batch of items to be processed // - cookie to be commited when processing is done // - error Next() (items []any, cookie int, err error) // Commit is used to mark data batch as processed Commit(cookie int) error}
type Consumer interface { Process(items []any) error}
func Pipe(p Producer, c Consumer) error { var buf []any var cookies []int for { items, cookie, err := p.Next() if err != nil { return fmt.Errorf("Producer errors: %w ", err) }
if len(buf)+len(items) <= MaxItems { buf = append(buf, items...) cookies = append(cookies, cookie) continue }
err := c.Process(buf) if err != nil { return fmt.Errorf("Consumer errors: %w ", err) }
for i, v := range cookies { err := p.Commit(v)
if err != nil { return fmt.Errorf("Not commit: %w ", err) } } buf = []any buf = append(buf, items...)
cookies = []int cookies = append(cookies, cookie) } return nil}