Sobes.tech
Middle

/* Մեզ անհրաժեշտ է փոխանցել տվյալներ որոշ աղբյուրից որոշ սպառողի: Աղբյուրը փոքր խմբերով է տալիս տվյալներ (~ տասնյակ գրառումներ), իսկ սպառողը ավելի արդյունավետ է աշխատում մեծ խմբերով: Օրինակ՝ Kafka տիպի հերթերից տվյալների փոխանցումը Clickhouse բազա: Աղբյուր: - Կանխամտածված անսահմանափակ: - Աղբյուրը երբեք չի վերադարձնի ավելի քան MaxItems գրանցամատյան մեկ կանչով Next: - Միևնույն "սեսիայի" շրջանակում (մեկ Pipe ֆունկցիայի կանչում) աղբյուրը յուրաքանչյուր Next-ով նոր տվյալներ է վերադարձնում: - Սակայն, վերագործարկումից հետո աղբյուրը սկսում է նախորդ "հաստատված" դիրքից, որը նշված է cookie-ով: Այդ պատճառով, յուրաքանչյուր cookie արժեք, որը վերադարձվել է Next-ով, պետք է հաստատվի Commit-ով, և դա պետք է կատարվի նույն հերթականությամբ, ինչպես վերադարձվել է: Սպառող: - Չի կարող մշակել ավելի քան MaxItems միանգամից: Հիմնական մակարդակ: Պետք է իրականացնել func Pipe(p Producer, c Consumer) error ֆունկցիան՝ որոնք կկարդա տվյալները աղբյուրից, կկազմեն դրանք բուֆեր, որի չափը չի գերազանցի MaxItems-ը, և կպահեն սպառողում, հետո կստորագրեն առաջընթացը աղբյուրում: */ const MaxItems = 9999 type Producer interface { // Next վերադարձնում է՝ // - խմբաքանակ տվյալների // - cookie հաստատելու համար // - սխալ Next() (items []any, cookie int, err error) // Commit-ը նշում է, որ տվյալ խմբաքանակը մշակված է 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

Պատասխան 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 в порядке получения, что соответствует требованиям.