Конкурентность и параллельная обработка
Все операции CyberGo JSON потокобезопасны, а параллельные API (ParallelIterator, параллельные JSONL-потоки) доступны из коробки. На этой странице — семантика потокобезопасности, встроенные параллельные API и паттерны конкурентного использования.
Подсказка Разделение со страницей производительности
Раздел «Конкурентная обработка» на странице Производительность показывает универсальные паттерны Go (sync.WaitGroup + семафор + Worker Pool) для ручного распараллеливания массивов; эта страница документирует встроенные в библиотеку параллельные API. Они дополняют друг друга.
Гарантии потокобезопасности
Processor — потокобезопасный движок обработки (комментарий в исходниках: Processor is the main JSON processing engine with thread safety):
- Один экземпляр Processor можно разделять между горутинами — все публичные методы (
Get/Set/Delete/Marshalи т.д.) внутренне защищены атомарными операциями и управлением конкурентности (beginGovernedOp/endGovernedOp). - Пакетные функции (
json.Get,json.GetStringи т.д.) разделяют один глобальный Processor и потокобезопасны по умолчанию. *ParsedJSONизPreParseможно читать конкурентно — несколько горутин могут одновременно вызыватьGetFromParsedдля одногоParsedJSON.
Предупреждение Когда не разделять
Processor можно разделять, но не разделяйте изменяемые Go-контейнеры между горутинами (например, передавать map[string]any из Get в несколько горутин для изменения). Возвращённые контейнеры по умолчанию — копии (если не включён CacheSharedResults), поэтому изменение возвращаемого значения не влияет на кэш. Но конкурентное изменение одного контейнера всё равно требует блокировки на стороне вызывающего.
ParallelIterator — параллельный итератор
ParallelIterator использует многоядерность CPU для параллельной обработки массивов: встроенный пул воркеров, агрегация ошибок и восстановление после panic — безопаснее самописного пула горутин.
Базовый параллельный обход
package main
import (
"fmt"
"sync"
"github.com/cybergodev/json"
)
func main() {
data := `{"items":[1,2,3,4,5,6,7,8]}`
items := json.GetArray(data, "items")
// Число воркеров по умолчанию = Config.MaxConcurrency (ограничено длиной массива)
iter := json.NewParallelIterator(items)
defer iter.Close()
var mu sync.Mutex
var sum int64
err := iter.ForEach(func(_ int, val any) error {
mu.Lock()
sum += int64(val.(float64))
mu.Unlock()
return nil
})
if err != nil {
panic(err)
}
fmt.Printf("Сумма = %d\n", sum)
// Вывод: Сумма = 36
}Параллельный Map
Map параллельно преобразует каждый элемент; результаты сохраняют порядок ввода (каждый воркер пишет в свой индекс, блокировка не нужна).
package main
import (
"fmt"
"github.com/cybergodev/json"
)
func main() {
data := `{"items":[1,2,3,4]}`
items := json.GetArray(data, "items")
iter := json.NewParallelIterator(items)
defer iter.Close()
// Параллельный map: каждый элемент * 10, порядок результата совпадает с вводом
doubled, err := iter.Map(func(_ int, val any) (any, error) {
return int(val.(float64)) * 10, nil
})
if err != nil {
panic(err)
}
fmt.Println(doubled)
// Вывод: [10 20 30 40]
}Обработка пакетами ForEachBatch / ForEachBatchWithContext
Когда накладные расходы колбэка на один элемент высоки (например, системный вызов или сетевой запрос на элемент), ForEachBatch нарезает элементы на пакеты фиксированного размера — каждый пакет обрабатывает одна горутина: внутри пакета последовательно, между пакетами параллельно, что распределяет стоимость планирования и синхронизации.
package main
import (
"context"
"fmt"
"time"
"github.com/cybergodev/json"
)
func main() {
data := `{"records":[10,20,30,40,50,60,70,80,90,100]}`
records := json.GetArray(data, "records")
iter := json.NewParallelIterator(records)
defer iter.Close()
// 10 записей по 3 на пакет -> 4 пакета (в последнем 1 запись); пишем по batchIdx в независимые индексы, блокировка не нужна
subtotals := make([]int, 4)
err := iter.ForEachBatch(3, func(batchIdx int, batch []any) error {
sum := 0
for _, v := range batch {
sum += int(v.(float64))
}
subtotals[batchIdx] = sum
return nil
})
if err != nil {
panic(err)
}
// Последовательное потребление после завершения всех пакетов (порядок выполнения не гарантируется, результаты расставлены по индексам)
for i, s := range subtotals {
fmt.Printf("Пакет %d: промежуточная сумма = %d\n", i, s)
}
// Версия с таймаутом: нераспределённые пакеты не запускаются после истечения ctx, запущенные выходят по отмене
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
defer cancel()
err = iter.ForEachBatchWithContext(ctx, 100, func(batchIdx int, batch []any) error {
return nil // имитация обработки одного пакета
})
fmt.Println("Пакетная обработка с таймаутом завершена, ошибка:", err)
}
// Вывод:
// Пакет 0: промежуточная сумма = 60
// Пакет 1: промежуточная сумма = 150
// Пакет 2: промежуточная сумма = 240
// Пакет 3: промежуточная сумма = 100
// Пакетная обработка с таймаутом завершена, ошибка: <nil>При batchSize <= 0 используется 100. Семантика ошибок колбэка как у ForEach: побеждает первая ошибка, диспетчеризация новых пакетов прекращается; порядок диспетчеризации пакетов совпадает со вводом (batchIdx растёт), но порядок выполнения не гарантируется — упорядоченный вывод достигается как в примере выше: расстановка по индексам с последовательным потреблением по завершении.
Обзор API ParallelIterator
| API | Сигнатура | Описание |
|---|---|---|
NewParallelIterator | func NewParallelIterator(data []any, cfg ...Config) *ParallelIterator | Создаёт итератор; число воркеров из cfg.MaxConcurrency (по умолчанию 50, ограничивается длиной массива; <= 0 откатывается к 4) |
ForEach | func (it *ParallelIterator) ForEach(fn func(int, any) error) error | Параллельный обход; возвращает первую ошибку |
ForEachWithContext | func (it *ParallelIterator) ForEachWithContext(ctx context.Context, fn func(int, any) error) error | Поддержка отмены через context |
ForEachBatch | func (it *ParallelIterator) ForEachBatch(batchSize int, fn func(int, []any) error) error | Параллельная обработка пакетами: внутри пакета последовательно, между пакетами параллельно |
ForEachBatchWithContext | func (it *ParallelIterator) ForEachBatchWithContext(ctx context.Context, batchSize int, fn func(int, []any) error) error | Пакетная параллельная обработка + отмена через context |
Map | func (it *ParallelIterator) Map(transform func(int, any) (any, error)) ([]any, error) | Параллельное преобразование с сохранением порядка |
Filter | func (it *ParallelIterator) Filter(predicate func(int, any) bool) []any | Параллельная фильтрация с сохранением порядка (без возврата ошибки) |
Close | func (it *ParallelIterator) Close() | Освобождает ресурсы (сигнализирует работающим горутинам остановиться; безопасно вызывать многократно) |
Полные сигнатуры и использование — в Типы итераторов.
Подсказка Обработка ошибок и panic
ForEach прекращает диспетчеризацию новых задач и возвращает первую ошибку; panic внутри воркера перехватывается (recover) и преобразуется в ошибку, поэтому panic в колбэке не роняет процесс. Для отмены используйте ForEachWithContext — корректно завершается при ctx.Done().
Параллельная потоковая обработка JSONL
Для больших JSONL (NDJSON) файлов StreamJSONLParallel обрабатывает каждую строку несколькими воркерами.
package main
import (
"fmt"
"strings"
"sync"
"github.com/cybergodev/json"
)
func main() {
// Имитация JSONL-данных (по одному JSON-объекту на строку)
jsonlData := `{"id":1,"score":95}
{"id":2,"score":82}
{"id":3,"score":78}
{"id":4,"score":90}`
processor, err := json.New()
if err != nil {
panic(err)
}
defer processor.Close()
var mu sync.Mutex
var total int64
var count int64
// 4 воркера обрабатывают строки параллельно
err = processor.StreamJSONLParallel(strings.NewReader(jsonlData), 4, func(lineNum int, item *json.IterableValue) error {
score := int64(item.GetInt("score"))
mu.Lock()
total += score
count++
mu.Unlock()
return nil
})
if err != nil {
panic(err)
}
fmt.Printf("Обработано %d записей, сумма %d\n", count, total)
// Вывод: Обработано 4 записей, сумма 345
}| API | Описание |
|---|---|
StreamJSONLParallel(reader, workers, fn) | Многоворкерная параллельная обработка JSONL |
StreamJSONLParallelWithContext(ctx, reader, workers, fn) | То же с отменой/таймаутом через context |
StreamJSONLChunked(reader, chunkSize, fn) | Обработка чанками, экономия памяти |
Полные сигнатуры и настройки (JSONLWorkers/JSONLChunkSize и т.д.) — в Обработка JSONL и Потоковая обработка JSONL.
Подсказка Порядок строк
В параллельном режиме lineNum в колбэке всё ещё отражает исходный номер строки, но порядок выполнения не гарантируется. Для упорядоченного вывода пишите в предвыделённый слайс по позиции lineNum.
Глобальный процессор для конкурентного использования
SetGlobalProcessor позволяет всем пакетным функциям разделять один пользовательский Processor — подходит для многогорутинных сервисов, требующих единой конфигурации (параметры кэша, хуки, ограничения безопасности).
package main
import (
"fmt"
"sync"
"github.com/cybergodev/json"
)
func main() {
// Пользовательский глобальный процессор (разделяется всеми пакетными функциями, потокобезопасен)
cfg := json.DefaultConfig()
processor, err := json.New(cfg)
if err != nil {
panic(err)
}
json.SetGlobalProcessor(processor) // предыдущий глобальный Processor закрывается автоматически
defer json.ShutdownGlobalProcessor() // корректное завершение при выходе
data := `{"user":{"name":"Alice","age":30}}`
// Несколько горутин конкурентно используют пакетные функции (один глобальный Processor)
var wg sync.WaitGroup
results := make([]string, 3)
for i := 0; i < 3; i++ {
wg.Add(1)
go func(idx int) {
defer wg.Done()
switch idx {
case 0:
results[idx] = json.GetString(data, "user.name")
case 1:
results[idx] = fmt.Sprintf("%d", json.GetInt(data, "user.age"))
case 2:
results[idx] = json.GetString(data, "user.name")
}
}(i)
}
wg.Wait()
fmt.Println(results)
// Вывод: [Alice 30 Alice]
}Предупреждение Передача владения
После SetGlobalProcessor жизненным циклом этого Processor управляет глобаль — не вызывайте на нём Close() вручную, иначе возникнет конфликт с логикой глобального завершения. Для корректного закрытия и освобождения ресурсов вызывайте ShutdownGlobalProcessor() при выходе.
Ограничение конкурентности MaxConcurrency
Config.MaxConcurrency (по умолчанию 50) — мягкий предел конкурентности на Processor: атомарный счётный семафор ограничивает число операций в полёте. При достижении предела новые операции возвращают ErrConcurrencyLimit (можно повторить).
cfg := json.DefaultConfig()
cfg.MaxConcurrency = 100 // повысить предел конкурентности на ProcessorErrConcurrencyLimit— повторяемая временная ошибка (см. Обработка ошибок).- Число воркеров параллельного стриминга (
StreamJSONLParallel) задаётся явным аргументом и не привязано напрямую кMaxConcurrency, но разделяет те же слоты управления. - Число воркеров
ParallelIteratorберётся изcfg.MaxConcurrency(по умолчанию 50), но ограничивается длиной массива.
Лучшие практики и подводные камни
1. Переиспользуйте Processor; не создавайте по одному на запрос
Processor хранит кэш, рекурсивный процессор и другое состояние — переиспользование одного экземпляра обеспечивает попадания кэша. Вызов json.New() на каждый запрос лишает преимуществ кэша и увеличивает аллокации.
2. Разделять экземпляр безопасно; разделять контейнеры результатов — осторожно
Processor безопасно разделять между горутинами; но map/slice, возвращённые Get, при совместном изменении между горутинами требуют блокировки на стороне вызывающего (или рассматривайте как только для чтения при включённом CacheSharedResults).
3. Освобождайте ресурсы через Close
В долгоживущих сервисах явно вызывайте defer processor.Close() и defer iter.Close(), чтобы избежать утечек горутин кэша и памяти. Экземпляр, установленный через SetGlobalProcessor, использует вместо этого ShutdownGlobalProcessor.
4. Параллелить стоит только CPU-интенсивную работу
У параллелизма есть накладные расходы на планирование и синхронизацию. Малые массивы (< ParallelThreshold, по умолчанию 10) быстрее обрабатываются последовательно; JSONL с большим числом строк и тяжёлой обработкой строки явно выигрывает от параллелизма.
5. Следите за порядком строк в параллельном режиме
StreamJSONLParallel не гарантирует порядок обработки. Для упорядоченного результата пишите по lineNum в позицию, затем потребляйте по порядку.
См. также
- Производительность — переиспользование Processor, универсальные паттерны конкурентности Go, бенчмарки
- Типы итераторов — полный API
ParallelIterator - Обработка JSONL — детали параллельного JSONL API
- Кэш и предпарсинг — механизм кэша и предпарсинг PreParse
- Обработка ошибок —
ErrConcurrencyLimitи классификация ошибок