/*
Биз манбаадан маълумотларни бир истеъмолчига ўтказишимиз керак.
Манба, маълумотларни кичик қисмлар (~ўн ёзув) шаклида беради, холос, истеъмолчи катта қисмлар (~минг ёзув) билан ишлашни самаралироқ қилади.
Ҳақиқий мисол — Kafka туридаги навбатлардан маълумотларни Clickhouse маълумотлар базасига ўтказиш.
Манба:
- Дейдарли чексиз.
- Манба ҳеч қачон Next чақирувида MaxItems дан ортиқ ёзув қайтармайди.
- Бир "сесия" (бир Pipe функцияси чақируви) давомида, манба ҳар бир Next чақирувида янги маълумотлар қайтариб туради.
- Аммо, қайта ишга туширгандан кейин, манба аввалги "тасдиқланган" позициядан, cookie билан белгиланган жойдан бошлайди.
Шунинг учун, *ҳар бир* cookie қиймати, Next қайтарган, маълумотлар сақлангандан кейин,
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 {
// TODO
}