Senior
/* Biz müəyyən bir mənbədən müəyyən bir istehlakçıya məlumat ötürməliyik. Mənbə kiçik partiyalar (~on qeyd) şəklində məlumat verir, halbuki, istehlakçı böyük partiyalar (~minlərlə qeyd) ilə daha səmərəli işləyir. Həqiqi nümunə Kafka tipli növbələrdən məlumatların Clickhouse bazasına ötürülməsidir. Mənbə: - Demək olar ki, sonsuz. - Mənbə heç vaxt Next çağırışında MaxItems-dən çox qeyd qaytarmır. - Bir "sessiya" (bir Pipe funksiyasının çağırışı) çərçivəsində, mənbə hər Next çağırışında yeni məlumatlar qaytarır. - Ancaq, yenidən başlatdıqdan sonra, mənbə əvvəlki "təsdiqlənmiş" mövqedən, cookie ilə göstərilən yerdən başlayır. Buna görə, Next tərəfindən qaytarılan hər cookie dəyəri, məlumatlar qəbul ediciyə saxlandıqdan sonra, Commit çağırışı ilə təsdiqlənməlidir, və yaşıdları ilə eyni ardıcıllıqda. Qəbul edən: - Eyni zamanda MaxItems-dən çox işləyə bilməz. Əsas səviyyə: func Pipe(p Producer, c Consumer) error funksiyasını həyata keçirmək lazımdır, bu mənbədən məlumat oxuyur, onları MaxItems-dən çox olmayan bufferdə qruplaşdırır və qəbulçiyə saxlayır, sonra isə mənbədə irəliləyişi təsdiqləyir. Çətinlik: Next, Process və Commit metodları şəbəkə çağırışları ilə əlaqəlidir və çox vaxt ala bilər. Prosesin sürətini artırmaq üçün, oxuma, yazma və irəliləyişi təsdiqləmə proseslərini paralel həyata keçirmək lazımdır. Belə ki, Process və ya Commit çağırışları zamanı, mənbədən oxuma və yeni bufferin formalaşması davam etsin. */ const MaxItems = 9999 type Producer interface { // Next aşağıdakıları qaytarır: // - işlənəcək elementlər partiyası // - təsdiqləmə üçün cookie // - xəta Next() (items []any, cookie int, err error) // Commit, məlumat partiyasını işlənmiş kimi işarələmək üçün istifadə olunur Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }
sobes.tech Süni İntellekt
AI-dan cavab
Для реализации функции 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, что позволяет прервать работу при возникновении ошибки.
Такой подход позволяет параллельно читать, обрабатывать и подтверждать данные, повышая общую производительность.