Middle
/* Szükségünk van adatok átadására egyik forrástól a másikhoz. A forrás kis csomagokban (~tíz rekord) adja az adatokat, míg a fogyasztó hatékonyabban működik nagyobb csomagokkal. Egy valós példa a Kafka típusú sorokból történő adatátvitel a Clickhouse adatbázisba. Forrás: - Szinte végtelen. - A forrás soha nem ad vissza több mint MaxItems rekordot egy Next hívásban. - Egy "munkamenet" (egy Pipe függvényhívás) során a forrás minden Next hívásnál új adatokat ad vissza. - Azonban újraindítás után a forrás a korábbi "megerősített" pozícióból kezd, amit cookie jelez. Ezért minden cookie érték, amit a Next visszaad, az adatok tárolása után, Commit hívással kell megerősíteni, ugyanabban a sorrendben, ahogyan a Next visszaadta. Fogyasztó: - Egyszerre nem dolgozhat több mint MaxItems-t. Alapszint: Implementálni kell a func Pipe(p Producer, c Consumer) error függvényt, ami adatokat olvas a forrásból, azokat MaxItems méretű bufferbe csoportosítja és a fogyasztóba menti, majd a forrásban a haladást megerősíti. */ const MaxItems = 9999 type Producer interface { // Next visszaad: // - egy csomagot az adatokból // - egy cookie-t, amit a feldolgozás végén kell megerősíteni // - hibát Next() (items []any, cookie int, err error) // Commit a feldolgozott csomag megerősítésére szolgál 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 err } buf = append(buf, items...) cookies = append(cookies, cookie) if len(buf) >= MaxItems { if err := c.Process(buf); err != nil { return err } for _, c := range cookies { if err := p.Commit(c); err != nil { return err } } buf = buf[:0] cookies = nil } } if len(buf) > 0 { if err := c.Process(buf); err != nil { return err } for _, c := range cookies { if err := p.Commit(c); err != nil { return err } } } return nil }
sobes.tech MI
Válasz az MI-től
Ваша задача — реализовать функцию Pipe, которая читает данные из источника (Producer) небольшими пачками, буферизует их до размера не более MaxItems, передает в Consumer и затем фиксирует прогресс в Producer через Commit в правильном порядке.
Основные моменты:
- Источник возвращает данные пачками, максимум MaxItems за вызов Next.
- Нужно аккумулировать данные в буфере, не превышая MaxItems.
- После передачи данных в Consumer вызываем Commit для всех cookie в порядке получения.
- Источник бесконечен, поэтому цикл бесконечный.
В представленном коде есть ошибка: бесконечный цикл не прерывается, и код после цикла недостижим. Нужно предусмотреть условие выхода (например, когда Next возвращает пустой срез и нет ошибки).
Исправленный пример реализации:
func Pipe(p Producer, c Consumer) error {
var buf []any
var cookies []int
for {
items, cookie, err := p.Next()
if err != nil {
return err
}
// Если источник вернул пустой срез, считаем, что данных больше нет
if len(items) == 0 {
break
}
buf = append(buf, items...)
cookies = append(cookies, cookie)
// Если буфер достиг максимума, отправляем на обработку
if len(buf) >= MaxItems {
if err := c.Process(buf); err != nil {
return err
}
for _, c := range cookies {
if err := p.Commit(c); err != nil {
return err
}
}
buf = buf[:0]
cookies = nil
}
}
// Обработка оставшихся данных
if len(buf) > 0 {
if err := c.Process(buf); err != nil {
return err
}
for _, c := range cookies {
if err := p.Commit(c); err != nil {
return err
}
}
}
return nil
}
Таким образом, мы аккумулируем данные, не превышая MaxItems, и после обработки фиксируем cookie в порядке получения, что соответствует требованиям.