Sobes.tech
Senior

/* We moeten gegevens van een bron naar een consument overbrengen. De bron geeft gegevens in kleine batches (~tien records), terwijl de consument 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 echter vanaf de vorige "bevestigde" positie, aangegeven door cookie. Daarom moet elke cookie-waarde die door Next wordt geretourneerd, na het opslaan van de gegevens in de ontvanger, worden bevestigd met een Commit-aanroep, in dezelfde volgorde als ze door Next werden geretourneerd. Ontvanger: - Kan niet meer dan MaxItems tegelijk verwerken. Basissysteem: 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. Uitdaging: De methoden Next, Process en Commit zijn verbonden met netwerkoproepen en kunnen vrij lang duren. Om het proces te versnellen, moeten de lees-, schrijf- en bevestigingsprocessen parallel worden uitgevoerd. Zodat tijdens Process of Commit, het lezen uit de bron en het vormen van de nieuwe buffer doorgaan. */ const MaxItems = 9999 type Producer interface { // Next geeft terug: // - een batch van items om te verwerken // - een cookie om te bevestigen wanneer de verwerking is voltooid // - een fout Next() (items []any, cookie int, err error) // Commit wordt gebruikt om een gegevensbatch als verwerkt te markeren Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }

sobes.tech AI

Antwoord van 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, что позволяет прервать работу при возникновении ошибки.

Такой подход позволяет параллельно читать, обрабатывать и подтверждать данные, повышая общую производительность.