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 arbeitet. Ein echtes Beispiel ist die Übertragung von Daten aus Kafka-Queues in die Clickhouse-Datenbank. Quelle: - Praktisch unendlich. - Die Quelle gibt nie mehr als MaxItems Datensätze in einem einzelnen Next-Aufruf zurück. - Innerhalb einer "Sitzung" (einem Aufruf der Funktion Pipe) liefert die Quelle bei jedem Next neue Daten. - Nach einem Neustart beginnt die Quelle wieder bei der vorherigen "bestätigten" Position, die durch cookie bestimmt wird. Daher muss jeder Wert von cookie, den Next zurückgibt, nach der Speicherung der Daten im Empfänger mit einem 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 mit maximaler Größe MaxItems gruppiert und 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 // - eine cookie, die nach Abschluss der Verarbeitung bestätigt werden soll // - einen 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 { var buf []any var cookies []int for { items, cookie, err := p.Next() if err != nil { return err } buf = append(buf, items...) cookies = append(cookies, cookie) if len(buf) >= MaxItems { if err := c.Process(buf); err != nil { return err } for _, c := range cookies { if err := p.Commit(c); err != nil { return err } } buf = buf[:0] cookies = nil } } if len(buf) > 0 { if err := c.Process(buf); err != nil { return err } for _, c := range cookies { if err := p.Commit(c); err != nil { return err } } } return nil }

sobes.tech KI

Antwort von AI

Ваша задача — реализовать функцию Pipe, которая читает данные из источника (Producer) небольшими пачками, буферизует их до размера не более MaxItems, передает в Consumer и затем фиксирует прогресс в Producer через Commit в правильном порядке.

Основные моменты:

  • Источник возвращает данные пачками, максимум MaxItems за вызов Next.
  • Нужно аккумулировать данные в буфере, не превышая MaxItems.
  • После передачи данных в Consumer вызываем Commit для всех cookie в порядке получения.
  • Источник бесконечен, поэтому цикл бесконечный.

В представленном коде есть ошибка: бесконечный цикл не прерывается, и код после цикла недостижим. Нужно предусмотреть условие выхода (например, когда Next возвращает пустой срез и нет ошибки).

Исправленный пример реализации:

func Pipe(p Producer, c Consumer) error {
    var buf []any
    var cookies []int

    for {
        items, cookie, err := p.Next()
        if err != nil {
            return err
        }
        // Если источник вернул пустой срез, считаем, что данных больше нет
        if len(items) == 0 {
            break
        }

        buf = append(buf, items...)
        cookies = append(cookies, cookie)

        // Если буфер достиг максимума, отправляем на обработку
        if len(buf) >= MaxItems {
            if err := c.Process(buf); err != nil {
                return err
            }
            for _, c := range cookies {
                if err := p.Commit(c); err != nil {
                    return err
                }
            }
            buf = buf[:0]
            cookies = nil
        }
    }

    // Обработка оставшихся данных
    if len(buf) > 0 {
        if err := c.Process(buf); err != nil {
            return err
        }
        for _, c := range cookies {
            if err := p.Commit(c); err != nil {
                return err
            }
        }
    }

    return nil
}

Таким образом, мы аккумулируем данные, не превышая MaxItems, и после обработки фиксируем cookie в порядке получения, что соответствует требованиям.