Middle
/* We moeten gegevens van een bron naar een ontvanger overbrengen. De bron levert de gegevens in kleine pakketten (~tientallen records), terwijl de ontvanger efficiënter werkt met grotere batches. Een echt voorbeeld is het overzetten van gegevens van Kafka-achtige wachtrijen naar de Clickhouse-database. Bron: - Bijna oneindig. - De bron geeft nooit meer dan MaxItems records terug in één Next-aanroep. - Tijdens één "sessie" (één aanroep van de Pipe-functie) geeft de bron bij elke Next nieuwe gegevens. - Na een herstart begint de bron weer vanaf de vorige "bevestigde" positie, aangegeven door cookie. Daarom moet elke waarde van cookie die Next teruggeeft, na het opslaan van de gegevens in de ontvanger, worden bevestigd met een Commit-aanroep, in dezelfde volgorde als waarin ze door Next werden teruggegeven. Ontvanger: - Kan niet meer dan MaxItems tegelijk verwerken. Basisniveau: Het is nodig om de functie func Pipe(p Producer, c Consumer) error te implementeren, die gegevens uit de bron leest, ze groepeert in een buffer van niet meer dan MaxItems en opslaat in de ontvanger, en vervolgens de voortgang in de bron bevestigt. */ const MaxItems = 9999 type Producer interface { // Next geeft terug: // - een batch van items om te verwerken // - een cookie om te bevestigen wanneer de verwerking is voltooid // - een fout Next() (items []any, cookie int, err error) // Commit wordt gebruikt om de gegevensbatch als verwerkt te markeren 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 AI
Antwoord van AI
Ваша задача — реализовать функцию 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 в порядке получения, что соответствует требованиям.