Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions cmd/mackerel-plugin-jsonl/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ type Opt struct {
filterByte *[]byte
ignoreByte *[]byte
paths [][]string
flatPaths bool
duration float64
}

Expand Down
38 changes: 17 additions & 21 deletions cmd/mackerel-plugin-jsonl/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"time"

"github.com/monitoring-forge/followparser"
"github.com/stretchr/testify/require"
)

func generateJSONLFile(b testing.TB, dir, filename string, numLines int) error {
Expand Down Expand Up @@ -90,7 +91,7 @@ func initParserForTest(b testing.TB, tmpDir string) (*followparser.Parser, *Opt)
return fp, opt
}

func internalBenchmarkParse(b *testing.B, numLines int, doOutput bool) {
func internalBenchmarkParse(b *testing.B, numLines int, doOutput, useEachKey bool) {
tmpDir := b.TempDir()
prefix := "json"
logFileName := "json.log"
Expand All @@ -107,44 +108,39 @@ func internalBenchmarkParse(b *testing.B, numLines int, doOutput bool) {
for b.Loop() {
b.StopTimer()
err := resetFollowParserStateFile(b, tmpDir, logFileName, prefix)
if err != nil {
b.Fatalf("resetFollowParserStateFile failed: %v", err)
}
require.NoError(b, err, "resetFollowParserStateFile failed")
b.StartTimer()
fp, opt := initParserForTest(b, tmpDir)
if useEachKey {
opt.flatPaths = false
}
parsed, err := fp.Parse(
posFile,
logFile,
)
if err != nil {
b.Fatalf("Parse failed: %v", err)
}
require.NoError(b, err, "Parse failed")
if doOutput {
output := opt.output()
if output == "" {
b.Fatalf("output is empty")
}
require.NotEmpty(b, output, "output is empty")
}
b.StopTimer()
if parsed == nil {
b.Fatalf("Parse returned nil parsed data")
}
if len(parsed) != 1 {
b.Fatalf("Parse returned unexpected number of parsed data: got %d, want 1", len(parsed))
}
if parsed[0].Rows != numLines {
b.Fatalf("Parse returned unexpected number of rows: got %d, want %d", parsed[0].Rows, numLines)
}
require.NotNil(b, parsed, "parsed data is nil")
require.Equal(b, 1, len(parsed), "unexpected number of parsed data")
require.Equal(b, numLines, parsed[0].Rows, "unexpected number of rows in parsed data")
b.StartTimer()
}
}

// generate 100k JSONL file and parse benchmark
func BenchmarkMainParse_jsonl(b *testing.B) {
internalBenchmarkParse(b, 100_000, false)
internalBenchmarkParse(b, 100_000, false, false)
}

func BenchmarkMainParse_jsonl_eachkey(b *testing.B) {
internalBenchmarkParse(b, 100_000, false, true)
}

// generate 100k JSONL file and parse benchmark
func BenchmarkMainParse_parse_and_output(b *testing.B) {
internalBenchmarkParse(b, 100_000, true)
internalBenchmarkParse(b, 100_000, true, false)
}
40 changes: 39 additions & 1 deletion cmd/mackerel-plugin-jsonl/parser.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,14 @@ package main

import (
"bytes"
"errors"
"log"

"github.com/buger/jsonparser"
)

var errFlatPathsFound = errors.New("all flat JSON paths found")

func (opt *Opt) jsonParsed(idx int, value []byte, vt jsonparser.ValueType, err error) {
if err != nil {
log.Printf("error: %v", err)
Expand Down Expand Up @@ -37,10 +40,45 @@ func (opt *Opt) Parse(b []byte) error {
}
}

jsonparser.EachKey(b, opt.jsonParsed, opt.paths...)
if opt.flatPaths {
opt.parseFlat(b)
} else {
jsonparser.EachKey(b, opt.jsonParsed, opt.paths...)
}
return nil
}

// parseFlat reads only root-level keys. ObjectEach decodes escaped key names,
// so comparisons here use the same names as paths passed to EachKey.
func (opt *Opt) parseFlat(b []byte) {
var found uint64
remaining := len(opt.paths)
err := jsonparser.ObjectEach(b, func(key, value []byte, valueType jsonparser.ValueType, _ int) error {
for i, path := range opt.paths {
bit := uint64(1) << i
if found&bit != 0 || !bytes.Equal(key, []byte(path[0])) {
continue
}
found |= bit
remaining--
opt.jsonParsed(i, value, valueType, nil)
}
if remaining == 0 {
return errFlatPathsFound
}
return nil
})
if err != nil && err != errFlatPathsFound { //nolint:errorlint
// Keep EachKey's behavior for non-object and malformed input. Skip values
// already delivered by ObjectEach so aggregators never count them twice.
jsonparser.EachKey(b, func(i int, value []byte, valueType jsonparser.ValueType, err error) {
if i < 0 || found&(uint64(1)<<i) == 0 {
opt.jsonParsed(i, value, valueType, err)
}
}, opt.paths...)
}
}

func (opt *Opt) Finish(duration float64) {
opt.duration = duration
}
104 changes: 104 additions & 0 deletions cmd/mackerel-plugin-jsonl/parser_test.go
Original file line number Diff line number Diff line change
@@ -1,11 +1,115 @@
package main

import (
"reflect"
"testing"

"github.com/monitoring-forge/sampdo"
"github.com/stretchr/testify/require"
)

func TestParseFlatMatchesEachKey(t *testing.T) {
tests := []struct {
name string
paths [][]string
line string
}{
{
name: "root keys after nested values",
paths: [][]string{{"time"}, {"status"}, {"reqtime"}},
line: `{"extra":{"status":"nested"},"time":"now","status":"200","reqtime":0.5}`,
},
{
name: "escaped key and value",
paths: [][]string{{"status"}, {"a.b"}},
line: `{"sta\u0074us":"ok\nnext","a.b":"dotted"}`,
},
{
name: "duplicate keys use first value",
paths: [][]string{{"status"}, {"time"}},
line: `{"status":"200","status":"500","time":"now"}`,
},
{
name: "duplicate requested paths",
paths: [][]string{{"status"}, {"status"}},
line: `{"status":"200","status":"500"}`,
},
{
name: "missing and null values",
paths: [][]string{{"status"}, {"time"}, {"missing"}},
line: `{"status":null,"time":"now"}`,
},
{
name: "compound values",
paths: [][]string{{"status"}, {"time"}, {"reqtime"}},
line: `{"status":{"code":200},"time":[1,2],"reqtime":true}`,
},
{
name: "root array uses generic parser",
paths: [][]string{{"status"}, {"time"}},
line: `[{"status":"200"},{"time":"now"}]`,
},
}

for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
newOpt := func() *Opt {
opt := &Opt{}
for _, path := range tc.paths {
opt.aggregatorFunctions = append(opt.aggregatorFunctions, &AggregatorFunction{
jsonKey: path, aggregator: "group_by", groupBy: map[string]int{},
})
}
opt.setupPaths()
return opt
}

fast := newOpt()
require.True(t, fast.flatPaths, "expected flat path optimization")
reference := newOpt()
reference.flatPaths = false
for _, opt := range []*Opt{fast, reference} {
err := opt.Parse([]byte(tc.line))
require.NoError(t, err, "Parse failed")
}
for i := range tc.paths {
if !reflect.DeepEqual(fast.aggregatorFunctions[i].groupBy, reference.aggregatorFunctions[i].groupBy) {
t.Errorf("path %v: fast=%v, EachKey=%v", tc.paths[i], fast.aggregatorFunctions[i].groupBy, reference.aggregatorFunctions[i].groupBy)
}
}
})
}
}

func TestSetupPathsSelectsFlatOnlyForSmallRootKeySets(t *testing.T) {
for _, tc := range []struct {
name string
paths [][]string
want bool
}{
{"one root key", [][]string{{"status"}}, true},
{"nested key", [][]string{{"status", "code"}}, false},
{"array key", [][]string{{"items", "[0]"}}, false},
{"no keys", nil, false},
{"64 keys", make([][]string, 64), true},
{"65 keys", make([][]string, 65), false},
} {
t.Run(tc.name, func(t *testing.T) {
opt := &Opt{}
for _, path := range tc.paths {
if path == nil {
path = []string{"status"}
}
opt.aggregatorFunctions = append(opt.aggregatorFunctions, &AggregatorFunction{jsonKey: path})
}
opt.setupPaths()
if opt.flatPaths != tc.want {
t.Errorf("flatPaths=%v, want %v", opt.flatPaths, tc.want)
}
})
}
}

func TestParser_Parse(t *testing.T) {
opt := &Opt{
aggregatorFunctions: []*AggregatorFunction{
Expand Down
10 changes: 9 additions & 1 deletion cmd/mackerel-plugin-jsonl/reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,10 @@ func (af *AggregatorFunction) appendData(b []byte) error {
if err != nil {
return err
}
af.percentiles.Append(floatValue)
err = af.percentiles.Append(floatValue)
if err != nil {
return err
}
}

return nil
Expand Down Expand Up @@ -133,8 +136,13 @@ func (p *Opt) setupFilterBytes() {

func (p *Opt) setupPaths() {
paths := make([][]string, 0, len(p.aggregatorFunctions))
// parseFlat uses a uint64 to track paths already found in each line.
p.flatPaths = len(p.aggregatorFunctions) > 0 && len(p.aggregatorFunctions) <= 64
for _, af := range p.aggregatorFunctions {
paths = append(paths, af.jsonKey)
if len(af.jsonKey) != 1 {
p.flatPaths = false
}
}
p.paths = paths
}
Expand Down
Loading