Senior
/* 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 besser 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 nie mehr als MaxItems Datensätze bei einem Aufruf von Next zurück. - Innerhalb einer "Sitzung" (einem Aufruf der Funktion Pipe) liefert die Quelle bei jedem Aufruf von Next neue Daten. - Nach einem Neustart beginnt die Quelle jedoch an der vorherigen "bestätigten" Position, die durch cookie bestimmt wird. Daher muss jeder Wert von cookie, der von Next zurückgegeben wird, nach dem Speichern der Daten im Empfänger mit Commit bestätigt werden, und zwar in genau der Reihenfolge, in der sie von Next zurückgegeben wurden. Empfänger: - Kann nicht mehr als MaxItems auf einmal verarbeiten. Basisebene: Es ist erforderlich, die Funktion func Pipe(p Producer, c Consumer) error zu implementieren, die Daten aus der Quelle liest, sie in einen Puffer gruppiert, der nicht größer als MaxItems ist, und sie im Empfänger speichert, und anschließend den Fortschritt in der Quelle bestätigt. Schwierigkeit: Die Methoden Next, Process und Commit sind mit Netzwerkaufrufen verbunden und können ziemlich lange dauern. Um den Prozess zu beschleunigen, müssen die Prozesse Lesen, Schreiben und Fortschrittsbestätigung parallelisiert werden. Damit während Process oder Commit das Lesen aus der Quelle und die Bildung des neuen Puffers fortgesetzt werden. */ const MaxItems = 9999 type Producer interface { // Next gibt zurück: // - einen Batch von Items zur Verarbeitung // - eine cookie, die bestätigt werden muss, wenn die Verarbeitung abgeschlossen ist // - einen Fehler Next() (items []any, cookie int, err error) // Commit wird verwendet, um einen Datenbatch als verarbeitet zu markieren Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }
sobes.tech KI
Antwort von AI
Для реализации функции Pipe с параллельной обработкой чтения, записи и подтверждения прогресса можно использовать каналы и горутины. Идея:
- Одна горутина читает данные из источника (Producer.Next), формируя батчи не более MaxItems.
- Вторая горутина обрабатывает батчи (Consumer.Process).
- Третья горутина подтверждает прогресс (Producer.Commit) строго в порядке получения cookie.
Для сохранения порядка подтверждения можно использовать очередь или канал с буфером, где cookie будут храниться в порядке поступления.
Пример реализации:
func Pipe(p Producer, c Consumer) error {
type batch struct {
items []any
cookie int
}
batches := make(chan batch, 10) // Канал для передачи батчей на обработку
commitCh := make(chan int, 10) // Канал для передачи cookie на подтверждение
errCh := make(chan error, 3) // Канал для ошибок
done := make(chan struct{}) // Канал для остановки
// Чтение данных
go func() {
defer close(batches)
for {
items, cookie, err := p.Next()
if err != nil {
errCh <- err
return
}
if len(items) == 0 {
// Источник пуст, можно завершать
return
}
// Формируем батчи не более MaxItems
for start := 0; start < len(items); start += MaxItems {
end := start + MaxItems
if end > len(items) {
end = len(items)
}
batches <- batch{items: items[start:end], cookie: cookie}
}
}
}()
// Обработка данных
go func() {
defer close(commitCh)
for b := range batches {
if err := c.Process(b.items); err != nil {
errCh <- err
return
}
commitCh <- b.cookie
}
}()
// Подтверждение
go func() {
for cookie := range commitCh {
if err := p.Commit(cookie); err != nil {
errCh <- err
return
}
}
close(done)
}()
select {
case err := <-errCh:
return err
case <-done:
return nil
}
}
В этом примере:
- Чтение и формирование батчей происходит в первой горутине.
- Обработка батчей — во второй.
- Подтверждение cookie — в третьей.
Каналы обеспечивают поток данных и порядок подтверждения. Ошибки передаются через errCh, что позволяет прервать работу при возникновении ошибки.
Такой подход позволяет параллельно читать, обрабатывать и подтверждать данные, повышая общую производительность.