/*
Biz manbadan ma'lumotlarni iste'molchiga uzatishimiz kerak.
Manba, ma'lumotlarni kichik partiyalar (~o'n yozuv) shaklida taqdim etadi, holbuki, iste'molchi katta partiyalar (~ming yozuv) bilan ishlashni optimallashtiradi.
Haqiqiy misol - Kafka turidagi navbatlardan ma'lumotlarni Clickhouse bazasiga uzatish.
Manba:
- deyarli cheksiz.
- Manba har doim Next chaqiruvida MaxItems dan ko'p yozuvlarni qaytirmaydi.
- Bir "sessiya" (bir Pipe funksiyasi chaqiruvi) doirasida, manba har 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 qaytaradigan, ma'lumotlar saqlangandan so'ng,
Commit chaqiruvi bilan tasdiqlanishi kerak, va ular Next tomonidan qaytarilgan tartibda.
Qabul qiluvchi:
- Bir vaqtning o'zida MaxItems dan ko'p ishlov berolmaydi.
Kerakli: func Pipe(p Producer, c Consumer) error funksiyasini amalga oshirish,
manbadan ma'lumotlarni o'qib, ularni MaxItems dan oshmaydigan bufferga guruhlab, qabul qiluvchiga saqlash,
va keyin manbada progressni tasdiqlash.
*/
const MaxItems = 9999
type Producer interface {
// Next:
// - ishlash uchun partiya elementlar
// - tugallangach tasdiqlash uchun cookie
// - xato
Next() (items []any, cookie int, err error)
// Commit, ma'lumotlar partiyasini ishlangan deb belgilash uchun ishlatiladi
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}