Конвейер JSONL-данных
В этом руководстве показано, как построить конвейер JSONL-данных (JSON с разделением строками) в CyberGo JSON: потоковое чтение, преобразование полей, пакетное преобразование формата и обработка больших файлов.
Потоковое чтение и преобразование JSONL
Используйте обобщённый StreamLinesInto[T] для построчного чтения JSONL-потока и десериализации в структуру, преобразуйте поля в функции обратного вызова, затем записывайте обратно в формат JSONL через ToJSONLString.
package main
import (
"fmt"
"strings"
"github.com/cybergodev/json"
)
// LogEntry представляет одну строку JSON-журнала
type LogEntry struct {
Timestamp string `json:"timestamp"`
Level string `json:"level"`
Message string `json:"message"`
}
// EnrichedLog — преобразованная запись (переименованные поля и новая категория)
type EnrichedLog struct {
Timestamp string `json:"ts"`
Level string `json:"level"`
Message string `json:"msg"`
Category string `json:"category"`
}
func main() {
// Имитация JSONL-потока журналов (на практике может поступать из файла или сети)
jsonlStream := `{"timestamp":"2024-01-01T10:00:00Z","level":"INFO","message":"служба запущена"}
{"timestamp":"2024-01-01T10:00:05Z","level":"ERROR","message":"сбой подключения к БД"}
{"timestamp":"2024-01-01T10:00:10Z","level":"WARN","message":"время ответа превысило порог"}
{"timestamp":"2024-01-01T10:00:15Z","level":"INFO","message":"переподключение выполнено"}`
reader := strings.NewReader(jsonlStream)
// 1. Потоковое чтение и преобразование каждой строки журнала
var enriched []any
entries, err := json.StreamLinesInto[LogEntry](reader, func(lineNum int, entry LogEntry) error {
// Категоризация по уровню
category := "normal"
if entry.Level == "ERROR" {
category = "critical"
} else if entry.Level == "WARN" {
category = "warning"
}
enriched = append(enriched, EnrichedLog{
Timestamp: entry.Timestamp,
Level: entry.Level,
Message: entry.Message,
Category: category,
})
return nil
})
if err != nil {
panic(err)
}
// 2. Пакетное преобразование обратно в формат JSONL
output, err := json.ToJSONLString(enriched)
if err != nil {
panic(err)
}
fmt.Printf("Обработано строк журнала: %d\n", len(entries))
fmt.Print(output)
}
// Вывод:
// Обработано строк журнала: 4
// {"ts":"2024-01-01T10:00:00Z","level":"INFO","msg":"служба запущена","category":"normal"}
// {"ts":"2024-01-01T10:00:05Z","level":"ERROR","msg":"сбой подключения к БД","category":"critical"}
// {"ts":"2024-01-01T10:00:10Z","level":"WARN","msg":"время ответа превысило порог","category":"warning"}
// {"ts":"2024-01-01T10:00:15Z","level":"INFO","msg":"переподключение выполнено","category":"normal"}Обработка JSONL-файлов
NDJSONProcessor обрабатывает JSONL-файл построчно; функция обратного вызова получает map[string]any (удобно, когда набор полей не фиксирован). Результаты агрегируются через ToJSONL для пакетного преобразования в JSONL-байты.
package main
import (
"fmt"
"os"
"path/filepath"
"github.com/cybergodev/json"
)
func main() {
// Создаём временный JSONL-файл, чтобы пример работал автономно
tmpDir, err := os.MkdirTemp("", "cybergo-pipeline-*")
if err != nil {
panic(err)
}
defer os.RemoveAll(tmpDir)
jsonlPath := filepath.Join(tmpDir, "events.jsonl")
jsonData := `{"event":"login","user":"alice","ts":"2024-01-01T10:00:00Z"}
{"event":"logout","user":"alice","ts":"2024-01-01T11:00:00Z"}
{"event":"login","user":"bob","ts":"2024-01-01T12:00:00Z"}
{"event":"purchase","user":"bob","ts":"2024-01-01T12:30:00Z"}`
if err := os.WriteFile(jsonlPath, []byte(jsonData), 0644); err != nil {
panic(err)
}
// 1. Построчная обработка через NDJSONProcessor (каждая строка парсится в map[string]any)
processor := json.NewNDJSONProcessor()
loginCount := 0
err = processor.ProcessFile(jsonlPath, func(lineNum int, obj map[string]any) error {
event, _ := obj["event"].(string)
user, _ := obj["user"].(string)
fmt.Printf("Строка %d: %s от %s\n", lineNum, event, user)
if event == "login" {
loginCount++
}
return nil
})
if err != nil {
panic(err)
}
// 2. Преобразование агрегированных результатов в JSONL (пакетное преобразование формата)
summary := []any{
map[string]any{"metric": "logins", "count": loginCount},
map[string]any{"metric": "total_events", "count": 4},
}
jsonlBytes, err := json.ToJSONL(summary)
if err != nil {
panic(err)
}
fmt.Printf("Событий входа: %d\n", loginCount)
fmt.Printf("Агрегированный результат:\n%s", string(jsonlBytes))
}
// Вывод:
// Строка 1: login от alice
// Строка 2: logout от alice
// Строка 3: login от bob
// Строка 4: purchase от bob
// Событий входа: 2
// Агрегированный результат:
// {"metric":"logins","count":2}
// {"metric":"total_events","count":4}Параллельный конвейер: StreamJSONLParallel + запись через JSONLWriter
Когда строк много, а обработка одной строки тяжёлая (преобразование, валидация, обогащение), StreamJSONLParallel потребляет поток несколькими worker'ами параллельно; результаты собираются в исходном порядке строк и записываются обратно в JSONL через JSONLWriter.WriteRaw без повторного кодирования:
package main
import (
"bytes"
"fmt"
"slices"
"strings"
"sync"
"github.com/cybergodev/json"
)
func main() {
// Имитация потока событий (на практике — большой файл: замените strings.NewReader на *os.File из os.Open)
jsonlStream := `{"event":"login","user":"alice","ts":"10:00"}
{"event":"page_view","user":"alice","ts":"10:01"}
{"event":"login","user":"bob","ts":"10:02"}
{"event":"purchase","user":"bob","ts":"10:03"}
{"event":"login","user":"carol","ts":"10:04"}`
p, err := json.New()
if err != nil {
panic(err)
}
defer p.Close()
// 1. Параллельная фильтрация и преобразование: оставляем только события login, переписываем в структуру {user, at}.
// Обратный вызов выполняется параллельно в нескольких worker'ах: при записи в общий map нужна блокировка;
// сохраняем по lineNum и после завершения восстанавливаем исходный порядок
var mu sync.Mutex
logins := make(map[int][]byte)
err = p.StreamJSONLParallel(strings.NewReader(jsonlStream), 4, func(lineNum int, item *json.IterableValue) error {
if item.GetString("event") != "login" {
return nil // пропускаем нецелевые события; возврат item.Break() корректно останавливает весь поток
}
encoded, err := json.Marshal(map[string]any{
"user": item.GetString("user"),
"at": item.GetString("ts"),
})
if err != nil {
return err // возвращённая ошибка останавливает диспетчеризацию и передаётся наверх как есть
}
mu.Lock()
logins[lineNum] = encoded
mu.Unlock()
return nil
})
if err != nil {
panic(err)
}
// 2. Записываем результат в исходном порядке строк (WriteRaw пишет уже закодированную строку, лишь добавляя перевод строки)
lineNums := make([]int, 0, len(logins))
for n := range logins {
lineNums = append(lineNums, n)
}
slices.Sort(lineNums)
var out bytes.Buffer
writer := json.NewJSONLWriter(&out)
for _, n := range lineNums {
if err := writer.WriteRaw(logins[n]); err != nil {
panic(err)
}
}
fmt.Printf("Отфильтровано событий входа: %d (записано строк: %d)\n", len(logins), writer.Stats().LinesProcessed)
fmt.Print(out.String())
}
// Вывод:
// Отфильтровано событий входа: 3 (записано строк: 3)
// {"at":"10:00","user":"alice"}
// {"at":"10:02","user":"bob"}
// {"at":"10:04","user":"carol"}Ключевые моменты параллельного конвейера
- Порядок: порядок выполнения параллельных обратных вызовов не гарантируется, но
lineNumвсегда соответствует исходному номеру строки — собирайте результаты по номерам строк, сортируйте и записывайте, чтобы сохранить порядок. - Число worker'ов: задаётся явно вторым аргументом (в примере — 4); при необходимости таймаута или отмены используйте
StreamJSONLParallelWithContext(ctx, reader, workers, fn). - Пропускная способность: выигрыш над последовательным
StreamJSONLзависит от стоимости обработки одной строки — лёгкие обратные вызовы чистого извлечения дают небольшой прирост, тяжёлые обратные вызовы обогащения/валидации — заметный.
Потоковая обработка больших файлов JSON-массивов
Для одного большого файла JSON-массива (не JSONL) используйте ForeachFile для пообъектного обхода без загрузки всего файла в память за раз.
package main
import (
"fmt"
"os"
"path/filepath"
"github.com/cybergodev/json"
)
func main() {
tmpDir, err := os.MkdirTemp("", "cybergo-big-*")
if err != nil {
panic(err)
}
defer os.RemoveAll(tmpDir)
// Создаём большой файл JSON-массива (имитация крупного набора данных)
arrayPath := filepath.Join(tmpDir, "records.json")
records := []any{
map[string]any{"id": 1, "amount": 100, "currency": "USD"},
map[string]any{"id": 2, "amount": 250, "currency": "EUR"},
map[string]any{"id": 3, "amount": 80, "currency": "USD"},
map[string]any{"id": 4, "amount": 500, "currency": "GBP"},
map[string]any{"id": 5, "amount": 120, "currency": "USD"},
}
if err := json.SaveToFile(arrayPath, records); err != nil {
panic(err)
}
// Потоковый обход каждого элемента массива через ForeachFile
p, err := json.New()
if err != nil {
panic(err)
}
defer p.Close()
totalUSD := 0
err = p.ForeachFile(arrayPath, func(key any, item *json.IterableValue) error {
currency := item.GetString("currency")
amount := item.GetInt("amount")
if currency == "USD" {
totalUSD += amount
}
return nil // вернуть item.Break() для досрочной остановки
})
if err != nil {
panic(err)
}
fmt.Printf("Сумма USD: %d\n", totalUSD)
}
// Вывод: Сумма USD: 320Подсказка
- JSONL-файлы (по одному независимому JSON-объекту на строку): используйте
StreamLinesInto[T],NDJSONProcessorилиStreamJSONLFile. - Большие файлы JSON-массивов (один JSON-массив со множеством элементов): используйте
ForeachFileдля потоковой обработки без полной загрузки в память.
Следующие шаги
- JSONL-потоковая обработка — полное руководство по обработке JSONL
- Обработка больших файлов — подробно о потоковой обработке больших файлов
- Базовые примеры — базовые чтение/запись JSONL
- Шпаргалка — быстрый справочник по API