Middle
/* Biz ma'lumotlarni manbadan iste'molchiga o'tkazishimiz kerak. Manba kichik partiyalar (taxminan o'n yozuv) bilan ma'lumotlarni taqdim etadi, iste'molch esa katta partiyalar bilan samaraliroq ishlaydi. Haqiqiy misol - Kafka turidagi navbatlardan Clickhouse bazasiga ma'lumot uzatish. Manba: - deyarli cheksiz. - Manba har doim Next chaqiruvida MaxItems dan ko'p yozuvlarni qaytirmaydi. - Bir "sessiya" (bir Pipe funksiyasi chaqiruvi) davomida, manba har Next chaqiruvida yangi ma'lumotlarni qaytaradi. - Biroq, qayta ishga tushirilgandan so'ng, manba avvalgi "tasdiqlangan" pozitsiyadan, cookie bilan belgilangan joydan boshlanadi. Shuning uchun, har bir cookie qiymati, Next qaytarib bergandan so'ng, ma'lumotlar qabul qiluvchiga saqlangandan keyin, Commit chaqiruvi bilan tasdiqlanishi kerak, va ular Next tomonidan qaytarilgan tartibda bo'lishi shart. Qabul qiluvchi: - Bir vaqtning o'zida MaxItems dan ko'p ishlov berolmaydi. Asosiy daraja: Kerak bo'lsa, func Pipe(p Producer, c Consumer) error funksiyasini amalga oshiring, bu funktsiya manbadan ma'lumotlarni o'qiydi, ularni MaxItems o'lchamdagi bufferga guruhlaydi va qabul qiluvchiga saqlaydi, so'ngra manbada progressni tasdiqlaydi. */ const MaxItems = 9999 type Producer interface { // Next quyidagilarni qaytaradi: // - ishlov beriladigan elementlar partiyasi // - tasdiqlash uchun cookie // - xato Next() (items []any, cookie int, err error) // Commit, ma'lumotlar partiyasini ishlov berilgan deb belgilash uchun ishlatiladi 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
AIdan javob
Ваша задача — реализовать функцию 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 в порядке получения, что соответствует требованиям.