Sobes.tech
Senior

/* Bir kaynaktan bir tüketiciye veri aktarmamız gerekiyor. Kaynak, verileri küçük paketler (~on kayıt) halinde sağlar, oysa tüketici büyük paketlerle (~bin kayıt) daha verimli çalışır. Gerçek bir örnek, Kafka türü kuyruklardan alınan verilerin Clickhouse veritabanına aktarılmasıdır. Kaynak: - Neredeyse sonsuz. - Kaynak, bir Next çağrısında hiçbir zaman MaxItems'ten fazla kayıt döndürmez. - Bir "oturum" (bir Pipe fonksiyon çağrısı) içinde, kaynak her Next çağrısında yeni veriler döndürür. - Ancak, yeniden başlatıldıktan sonra, kaynak önceki "onaylanmış" konumdan, cookie ile belirlenen noktadan başlar. Bu nedenle, Next tarafından döndürülen her cookie değeri, veriler alıcıya kaydedildikten sonra, aynı sırayla Commit çağrısı ile onaylanmalıdır. Alıcı: - Aynı anda MaxItems'ten fazla işleyemez. Temel seviye: İşte, p Producer ve c Consumer fonksiyonlarını kullanarak, kaynaktan veri okuyan, bunları MaxItems boyutunu aşmayan bir tamponda toplayan ve alıcıya kaydeden, ve ardından kayıtlardaki ilerlemeyi onaylayan bir fonksiyon uygulaması gerekmektedir. Zorluk: Next, Process ve Commit metodları ağ çağrılarıyla ilgilidir ve oldukça uzun sürebilir. İşlemi hızlandırmak için, okuma, yazma ve ilerleme onay süreçlerini paralelleştirmek gerekir. Böylece, Process veya Commit çağrısı sırasında, kaynak okuma ve yeni tamponun oluşturulması devam eder. */ const MaxItems = 9999 type Producer interface { // Next şu değerleri döner: // - işlenmek üzere bir öğe grubu // - işleme tamamlandığında onaylanacak cookie // - hata Next() (items []any, cookie int, err error) // Commit, işlenmiş veri paketini işaretlemek için kullanılır Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }

sobes.tech yapay zeka

AI'dan gelen yanıt

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

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