Sobes.tech
Senior

/* Մեզ անհրաժեշտ է տվյալներ փոխանցել որոշ աղբյուրից որոշ սպառողի: Աղբյուրը տվյալները տալիս է փոքր խմբաքանակներով (~տասը գրառում), մինչդեռ սպառողը ավելի արդյունավետ է աշխատում մեծ խմբաքանակներով (~հազար գրառում): Իրական օրինակ է Kafka տիպի հերթերից տվյալների փոխանցումը Clickhouse տվյալաբազա: Աղբյուր: - Հասկանալիորեն անսահմանափակ: - Աղբյուրը երբեք չի վերադարձնի ավելի քան MaxItems գրանցում մեկ Next կանչում: - Մի «սեսիայում» (մեկ Pipe ֆունկցիայի կանչում) աղբյուրը յուրաքանչյուր Next կանչում նոր տվյալներ է վերադարձնում: - Սակայն, վերագործարկումից հետո, աղբյուրը սկսում է նախորդ «հաստատված» դիրքից, որը նշվում է cookie-ով: Այդ պատճառով, յուրաքանչյուր cookie արժեք, որը Next-ը վերադարձնում է, տվյալները պահելուց հետո, պետք է հաստատվի Commit կանչով, նույն հերթականությամբ, ինչպես դրանք վերադարձվել էին Next-ով: Սպառող: - Չի կարող միաժամանակ ավելի քան MaxItems մշակել: Հիմնական մակարդակ: Պետք է իրականացնել func Pipe(p Producer, c Consumer) error ֆունկցիան, որը կարդում է տվյալները աղբյուրից, խմբավորում դրանք ոչ մեծ, MaxItems-ից ոչ ավելի մեծ բուֆերներում և պահում է սպառողում, հետո հաստատում է առաջընթացը աղբյուրում: Ավելորդ բարդություն: Next, Process և Commit մեթոդները կապված են ցանցային կանչերի հետ և կարող են բավական երկար տևել: Արագացման համար անհրաժեշտ է հոսքերը հորիզոնականացնել՝ ընթերցում, գրառում և առաջընթացի հաստատում: Որպեսզի, Process կամ Commit կանչի ժամանակ, ընթերցումը աղբյուրից և նոր բուֆերի ձևավորումը շարունակվի: */ const MaxItems = 9999 type Producer interface { // Next վերադարձնում է: // - մի խմբաքանակ տարրեր՝ մշակման համար // - cookie՝ հաստատելու համար, երբ մշակումը ավարտվի // - սխալ Next() (items []any, cookie int, err error) // Commit օգտագործվում է նշելու տվյալների խմբաքանակը որպես մշակված 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, что позволяет прервать работу при возникновении ошибки.

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