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