Перейти к содержимому

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
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
}