/*
Mums jānodod dati no avota uz patērētāju.
Avots sniedz datus mazās partijās (~desmit ieraksti), kamēr patērētājs efektīvāk strādā ar lielākām partijām (~tūkstotis ierakstu).
Reāls piemērs ir datu pārnešana no Kafka tipa rindām uz Clickhouse datu bāzi.
Avots:
- Praktiski bezgalīgs.
- Avots nekad neatgriež vairāk par MaxItems ierakstiem vienā Next izsaukumā.
- Vienas "sesijas" (viena Pipe funkcijas izsaukuma) laikā avots katru reizi atgriež jaunus datus.
- Tomēr, pēc restartēšanas, avots sāk no iepriekšējās "apstiprinātās" pozīcijas, kas noteikta ar cookie.
Tādēļ, *katra* cookie vērtība, ko Next atgriež, pēc datu saglabāšanas saņēmējā,
ir jāapstiprina ar Commit izsaukumu, tādā pašā secībā, kādā tie tika atgriezti no Next.
Saņēmējs:
- Nevar apstrādāt vairāk par MaxItems vienlaikus.
Nepieciešams īstenot funkciju func Pipe(p Producer, c Consumer) error,
kura lasa datus no avota, tos grupē buferī ar izmēru ne lielāku par MaxItems un saglabā saņēmējā,
pēc tam apstiprina progresu avotā.
*/
const MaxItems = 9999
type Producer interface {
// Next atgriež:
// - partiju elementu, kurus apstrādāt
// - cookie, ko apstiprināt, kad apstrāde ir pabeigta
// - kļūdu
Next() (items []any, cookie int, err error)
// Commit tiek izmantots, lai marķētu datu partiju kā apstrādātu
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}