Sobes.tech
Senior

/* Nous devons transférer des données d'une source à un consommateur. La source fournit les données en petits lots (~dix enregistrements), tandis que le consommateur fonctionne mieux avec de gros lots (~mille enregistrements). Un exemple réel est le transfert de données depuis des files d'attente de type Kafka vers une base de données Clickhouse. Source : - Pratiquement infinie. - La source ne renvoie jamais plus de MaxItems enregistrements lors d'un appel à Next. - Dans le cadre d'une "session" (un appel à la fonction Pipe), la source renvoie de nouvelles données à chaque appel à Next. - Cependant, après un redémarrage, la source reprend à la position "confirmée" précédente, indiquée par cookie. Par conséquent, *chaque* valeur de cookie renvoyée par Next, après avoir sauvegardé les données dans le récepteur, doit être confirmée avec l'appel à Commit, dans le même ordre dans lequel elles ont été renvoyées par Next. Récepteur : - Ne peut traiter plus de MaxItems à la fois. Niveau de base : Il faut implémenter la fonction func Pipe(p Producer, c Consumer) error qui lit les données de la source, les regroupe dans un buffer de taille maximale MaxItems et les stocke dans le récepteur, après quoi elle confirme la progression dans la source. Difficulté supplémentaire : Les méthodes Next, Process et Commit sont liées à des appels réseau et peuvent prendre beaucoup de temps. Pour accélérer le processus, il faut paralléliser les processus de lecture, d'écriture et de confirmation de progression. De sorte qu'au cours de Process ou Commit, la lecture de la source et la formation du nouveau buffer continuent. */ const MaxItems = 9999 type Producer interface { // Next retourne : // - un lot d'éléments à traiter // - un cookie à confirmer une fois le traitement terminé // - une erreur Next() (items []any, cookie int, err error) // Commit est utilisé pour marquer un lot de données comme traité Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }

sobes.tech IA

Réponse de l'IA

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

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