Sobes.tech
Middle+

/* Wir müssen Daten von einer Quelle zu einem Verbraucher übertragen. Die Quelle liefert die Daten in kleinen Chargen (~zehn Einträge), während der Verbraucher effizienter mit größeren Chargen (~tausend Einträge) arbeitet. Ein echtes Beispiel ist die Übertragung von Daten aus Kafka-Queues in eine Clickhouse-Datenbank. Quelle: - Praktisch unendlich. - Die Quelle gibt niemals mehr als MaxItems Einträge in einem einzelnen Next-Aufruf zurück. - Innerhalb einer "Sitzung" (einen Aufruf der Funktion Pipe) liefert die Quelle bei jedem Next-Aufruf neue Daten. - Nach einem Neustart beginnt die Quelle wieder bei der vorherigen "bestätigten" Position, die durch cookie festgelegt ist. Daher muss jeder Wert von cookie, den Next zurückgibt, nach dem Speichern der Daten im Empfänger mit einem Commit bestätigt werden, und zwar in der gleichen Reihenfolge, in der sie von Next zurückgegeben wurden. Empfänger: - Kann nicht mehr als MaxItems auf einmal verarbeiten. Es ist erforderlich, die Funktion func Pipe(p Producer, c Consumer) error zu implementieren, die Daten aus der Quelle liest, sie in einen Puffer mit maximaler Größe MaxItems gruppiert und sie im Empfänger speichert, und anschließend den Fortschritt in der Quelle bestätigt. */ const MaxItems = 9999 type Producer interface { // Next gibt zurück: // - einen Batch von Items zur Verarbeitung // - cookie, das bei Abschluss der Verarbeitung bestätigt werden muss // - Fehler Next() (items []any, cookie int, err error) // Commit wird verwendet, um den Datenbatch als verarbeitet zu markieren Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }