Sobes.tech
Middle

/* Nous devons transférer des données d'une source à un consommateur. La source fournit les données en petits lots (~dizaines d'enregistrements), tandis que le consommateur fonctionne plus efficacement avec de gros lots. Un exemple réel est le transfert de données depuis des files d'attente de type Kafka vers la base de données Clickhouse. Source: - Pratiquement infinie. - La source ne renvoie jamais plus de MaxItems en un seul appel à Next. - Dans 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 recommencera à partir de 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 un appel à Commit, dans le même ordre dans lequel elles ont été renvoyées par Next. Récepteur: - Ne peut pas traiter plus de MaxItems à la fois. Niveau de base: Il est nécessaire d'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 sauvegarde dans le récepteur, puis confirme la progression dans la source. */ const MaxItems = 9999 type Producer interface { // Next retourne: // - un lot d'éléments à traiter // - un cookie à confirmer lorsque le traitement est terminé // - une erreur Next() (items []any, cookie int, err error) // Commit est utilisé pour marquer le lot de données comme traité 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 IA

Réponse de l'IA

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