Sobes.tech
Middle+

/* We moeten gegevens van een bron naar een ontvanger overbrengen. De bron geeft gegevens in kleine batches (~tien records), terwijl de ontvanger efficiënter werkt met grote batches (~duizend records). Een echt voorbeeld is het overzetten van gegevens van Kafka-achtige wachtrijen naar een Clickhouse-database. Bron: - Bijna oneindig. - De bron geeft nooit meer dan MaxItems records terug in één Next-aanroep. - Tijdens één "sessie" (één aanroep van de Pipe-functie) geeft de bron bij elke Next-aanroep nieuwe gegevens. - Na een herstart begint de bron vanaf de vorige "bevestigde" positie, bepaald door cookie. Daarom moet elke cookie-waarde die Next teruggeeft, na het opslaan van gegevens in de ontvanger, worden bevestigd met een Commit-aanroep, in dezelfde volgorde als waarin ze door Next werden teruggegeven. Ontvanger: - Kan niet meer dan MaxItems tegelijk verwerken. Het is nodig om de functie func Pipe(p Producer, c Consumer) error te implementeren, die gegevens uit de bron leest, deze groepeert in een buffer van niet groter dan MaxItems en opslaat in de ontvanger, waarna de voortgang in de bron wordt bevestigd. */ const MaxItems = 9999 type Producer interface { // Next geeft: // - een batch van items om te verwerken // - cookie om te bevestigen wanneer de verwerking is voltooid // - fout Next() (items []any, cookie int, err error) // Commit wordt gebruikt om de gegevensbatch als verwerkt te markeren Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }