Sobes.tech
Senior

/* Biz ma'lumotlarni manbadan iste'molchiga o'tkazishimiz kerak. Manba kichik partiyalar (~o'n yozuv) shaklida ma'lumotlarni taqdim etadi, holbuki, iste'molchi katta partiyalar (~ming yozuv) bilan ishlashni afzal ko'radi. Haqiqiy misol - Kafka turidagi navbatlardan ma'lumotlarni Clickhouse bazasiga uzatish. Manba: - deyarli cheksiz. - Manba har doim Next chaqiruvida MaxItems dan ko'p yozuvlarni qaytirmaydi. - Bir "sessiya" (bir Pipe funksiyasi chaqiruvi) davomida, manba har chaqiruvida yangi ma'lumotlarni qaytaradi. - Biroq, qayta ishga tushirilgandan so'ng, manba avvalgi "tasdiqlangan" pozitsiyadan, cookie bilan belgilangan joydan boshlanadi. Shuning uchun, Next tomonidan qaytarilgan har bir cookie qiymati, ma'lumotlar qabul qiluvchiga saqlangandan so'ng, Commit chaqiruvi bilan tasdiqlanishi kerak, vaqti-vaqti bilan Next tomonidan qaytarilgan tartibda. Qabul qiluvchi: - Bir vaqtning o'zida MaxItems dan ko'p ishlay olmaydi. Asosiy daraja: func Pipe(p Producer, c Consumer) error funksiyasini amalga oshirish talab etiladi, bu manbadan ma'lumotlarni o'qib, ularni MaxItems dan oshmaydigan bufferga to'playdi va qabul qiluvchiga saqlaydi, so'ngra manbada progressni tasdiqlaydi. Qiyinchilik: Next, Process va Commit metodlari tarmoq chaqiriqlari bilan bog'liq bo'lib, juda uzoq davom etishi mumkin. Jarayonni tezlashtirish uchun, o'qish, yozish va progressni tasdiqlash jarayonlarini parallel ravishda amalga oshirish kerak. Shunday qilib, Process yoki Commit davomida, manbadan o'qish va yangi bufferni shakllantirish davom etadi. */ const MaxItems = 9999 type Producer interface { // Next quyidagilarni qaytaradi: // - ishlash uchun elementlar partiyasi // - tasdiqlash uchun cookie // - xato Next() (items []any, cookie int, err error) // Commit, ma'lumotlar partiyasini ishlangan deb belgilash uchun ishlatiladi Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }

sobes.tech AI

AIdan javob

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

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