Sobes.tech
Middle+

/* 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 }