Sobes.tech
Middle+

/* Musimy przekazać dane z pewnego źródła do pewnego odbiorcy. Źródło zwraca dane w małych partiach (~dziesięć rekordów), podczas gdy odbiorca działa efektywniej z dużymi partiami (~tysiąc rekordów). Przykład rzeczywisty to przesyłanie danych z kolejek typu Kafka do bazy danych Clickhouse. Źródło: - Praktycznie nieskończone. - Źródło nigdy nie zwraca więcej niż MaxItems rekordów w jednym wywołaniu Next. - W ramach jednej "sesji" (jednego wywołania funkcji Pipe), źródło za każdym razem zwraca nowe dane na każde wywołanie Next. - Jednak po restarcie, źródło zacznie od poprzedniej "potwierdzonej" pozycji, określanej przez cookie. Dlatego *każda* wartość cookie, którą zwróci Next, po zapisaniu danych w odbiorniku, musi być potwierdzona wywołaniem Commit, w tej samej kolejności, w jakiej zostały zwrócone przez Next. Odbiornik: - Nie może przetworzyć więcej niż MaxItems naraz. Wymagana jest implementacja funkcji func Pipe(p Producer, c Consumer) error która czyta dane ze źródła, grupuje je w bufor o rozmiarze nie większym niż MaxItems i zapisuje do odbiornika, a następnie potwierdza postęp w źródle. */ const MaxItems = 9999 type Producer interface { // Next zwraca: // - partę elementów do przetworzenia // - cookie do potwierdzenia po zakończeniu przetwarzania // - błąd Next() (items []any, cookie int, err error) // Commit służy do oznaczenia partii danych jako przetworzonych Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }