Sobes.tech
Senior

/* Πρέπει να μεταφέρουμε δεδομένα από μια πηγή σε έναν καταναλωτή. Η πηγή παρέχει δεδομένα σε μικρές παρτίδες (~δέκα εγγραφές), ενώ ο καταναλωτής λειτουργεί πιο αποδοτικά με μεγάλες παρτίδες (~χίλιες εγγραφές). Ένα πραγματικό παράδειγμα είναι η μεταφορά δεδομένων από ουρές τύπου Kafka σε μια βάση δεδομένων Clickhouse. Πηγή: - Πρακτικά ατελείωτη. - Η πηγή ποτέ δεν επιστρέφει περισσότερα από MaxItems εγγραφές σε μια κλήση Next. - Σε μια "συνεδρία" (μία κλήση της λειτουργίας Pipe), η πηγή επιστρέφει νέα δεδομένα σε κάθε κλήση Next. - Ωστόσο, μετά από επανεκκίνηση, η πηγή ξεκινά από την προηγούμενη "επιβεβαιωμένη" θέση, που υποδεικνύεται από το cookie. Επομένως, κάθε τιμή cookie που επιστρέφεται από το Next, μετά την αποθήκευση των δεδομένων στον δέκτη, πρέπει να επιβεβαιωθεί με μια κλήση στο Commit, με την ίδια σειρά με αυτή που επιστράφηκε από το Next. Δέκτης: - Δεν μπορεί να επεξεργαστεί περισσότερα από MaxItems ταυτόχρονα. Βασικό επίπεδο: Απαιτείται η υλοποίηση της λειτουργίας func Pipe(p Producer, c Consumer) error που διαβάζει δεδομένα από την πηγή, τα ομαδοποιεί σε ένα buffer μεγέθους όχι μεγαλύτερο από MaxItems και τα αποθηκεύει στον δέκτη, μετά επιβεβαιώνει την πρόοδο στην πηγή. Πρόσθετη δυσκολία: Οι μέθοδοι Next, Process και Commit σχετίζονται με κλήσεις δικτύου και μπορεί να διαρκέσουν αρκετά. Για την επιτάχυνση της διαδικασίας, πρέπει να παραλληλοποιηθούν οι διαδικασίες ανάγνωσης, εγγραφής και επιβεβαίωσης προόδου. Έτσι, κατά τη διάρκεια των Process ή Commit, η ανάγνωση από την πηγή και η δημιουργία του νέου buffer συνεχίζονται. */ const MaxItems = 9999 type Producer interface { // Next επιστρέφει: // - ένα batch αντικειμένων προς επεξεργασία // - ένα cookie για επιβεβαίωση όταν η επεξεργασία ολοκληρωθεί // - ένα σφάλμα Next() (items []any, cookie int, err error) // Το Commit χρησιμοποιείται για να σηματοδοτήσει ένα batch δεδομένων ως επεξεργασμένο Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }

sobes.tech AI

Απάντηση από 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, что позволяет прервать работу при возникновении ошибки.

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