Sobes.tech
Middle

/* අපට මූලාශ්‍රයකින් දත්ත කිසියම් පරිභෝජකයකුට මාරු කළ යුතුය. මේදී මූලාශ්‍රය කුඩා කට්ටලවලින් (~ දස ගණනක් වාර්තා) දත්ත ලබා දී ඇති අතර, පරිභෝජකයා විශාල බචක සමඟ වැඩ කිරීම වඩාත් ප්‍රයෝජනවත් වේ. සත්‍ය උදාහරණයක් - Kafka වර්ගයේ පෙළගැස්මවලින් දත්ත ලබා දීම Clickhouse දත්තගබඩාවට. මූලාශ්‍රය: - කොන්දේසි ලෙස අසීමිත. - මූලාශ්‍රය එක් කැඳවීමකදී MaxItems වඩා වැඩි වාර්තා ආපසු නොදේ. - එක "සෙෂන්" (Pipe ක්‍රියාවලියක්) තුළ, මූලාශ්‍රය සෑම කැඳවීමකදී නව දත්ත ආපසු දක්වයි. - එහෙත්, නැවත ආරම්භ කිරීමෙන් පසු, මූලාශ්‍රය පසුගිය "තහවුරු කළ" තැනින් ආරම්භ වේ, එය cookie මගින් නියමිතයි. එමනිසා, Next කැඳවීමෙන් ආපසු ලැබූ සෑම cookie අගයක්ම, දත්ත සුරක්ෂිත කිරීමෙන් පසු, Commit මඟින් තහවුරු කළ යුතුය, එය ඒම අනුක්‍රමිකත්වයෙන්, Next විසින් ආපසු ලබා දී ඇත්තේ මෙන්. පරිභෝජක: - 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 в порядке получения, что соответствует требованиям.