Sobes.tech
Middle

/* Trebuie să transferăm date dintr-o sursă către un consumator. Sursa oferă date în loturi mici (~zeci de înregistrări), în timp ce consumatorul funcționează mai eficient cu loturi mari. Un exemplu real este transferul de date din cozi de tip Kafka în baza de date Clickhouse. Sursa: - Practic infinită. - Sursa nu returnează niciodată mai mult de MaxItems în cadrul unei singure apelări Next. - În cadrul unei "sesiuni" (un apel al funcției Pipe), sursa returnează date noi la fiecare Next. - Cu toate acestea, după repornire, sursa va începe de la poziția "confirmată" anterioară, indicată de cookie. Prin urmare, *fiecare* valoare de cookie returnată de Next, după salvarea datelor în receptor, trebuie confirmată cu o apelare la Commit, în aceeași ordine în care au fost returnate de Next. Receptor: - Nu poate procesa mai mult de MaxItems odată. Nivel de bază: Este necesar să implementați funcția func Pipe(p Producer, c Consumer) error care citește date din sursă, le grupează într-un buffer de dimensiune maximă MaxItems și le salvează în receptor, apoi confirmă progresul în sursă. */ const MaxItems = 9999 type Producer interface { // Next returnează: // - un lot de elemente de procesat // - un cookie pentru a fi confirmat după finalizarea procesării // - o eroare Next() (items []any, cookie int, err error) // Commit este folosit pentru a marca lotul de date ca fiind procesat 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

Răspuns de la 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 в порядке получения, что соответствует требованиям.