Sobes.tech
Middle+

/* Трябва да прехвърлим данни от източник към потребител. Източникът предоставя данни на малки пакети (~десет записа), докато потребителят работи по-ефективно с големи пакети (~хиляда записа). Реален пример е прехвърлянето на данни от опашки тип Kafka към база данни Clickhouse. Източник: - Практически безкраен. - Източникът никога не връща повече от MaxItems записи в един обаждане на Next. - В рамките на една "сесия" (едно обаждане на функцията 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 }