-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathprocessor_func_bench_test.go
More file actions
178 lines (163 loc) · 5.29 KB
/
Copy pathprocessor_func_bench_test.go
File metadata and controls
178 lines (163 loc) · 5.29 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
// Copyright © 2026 Meroxa, Inc.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package sdk
import (
"context"
"fmt"
"testing"
"github.com/conduitio/conduit-commons/opencdc"
)
// benchBatchSize is the number of records processed per Process call in the
// benchmarks below. It is representative of a typical pipeline batch, not a
// single record; ns/op and allocs/op should be divided by this constant to
// get a per-record figure.
const benchBatchSize = 100
func benchBatchRaw(n int) []opencdc.Record {
records := make([]opencdc.Record, n)
for i := range records {
records[i] = opencdc.Record{
Position: opencdc.Position("bench-position"),
Operation: opencdc.OperationUpdate,
Metadata: opencdc.Metadata{
"collection": "orders",
"version": "v1",
},
Key: opencdc.RawData("order-12345"),
Payload: opencdc.Change{
Before: opencdc.RawData(`{"id":12345,"status":"pending"}`),
After: opencdc.RawData(`{"id":12345,"status":"shipped"}`),
},
}
}
return records
}
func benchBatchStructured(n int) []opencdc.Record {
records := make([]opencdc.Record, n)
for i := range records {
records[i] = opencdc.Record{
Position: opencdc.Position("bench-position"),
Operation: opencdc.OperationUpdate,
Metadata: opencdc.Metadata{
"collection": "orders",
"version": "v1",
},
Key: opencdc.StructuredData{
"id": 12345,
},
Payload: opencdc.Change{
Before: opencdc.StructuredData{
"id": 12345,
"status": "pending",
"customer": map[string]any{
"id": 987,
"name": "Jane Doe",
},
},
After: opencdc.StructuredData{
"id": 12345,
"status": "shipped",
"customer": map[string]any{
"id": 987,
"name": "Jane Doe",
},
},
},
}
}
return records
}
// BenchmarkProcessorFunc_Process_MetadataFieldSet benchmarks a representative
// field-set style processor (built with the ProcessorFunc adapter, the SDK's
// concrete, reusable Processor implementation) that resolves a reference once
// per record and sets a metadata field on it. This is the same
// resolve-then-set pattern used by real field-set/field-rename processors,
// exercised through the full Processor.Process call including the
// ProcessedRecord wrapping. The reference works against both raw and
// structured payloads since it targets metadata, not the payload itself.
func BenchmarkProcessorFunc_Process_MetadataFieldSet(b *testing.B) {
resolver, err := NewReferenceResolver(".Metadata.processed_by")
if err != nil {
b.Fatal(err)
}
proc := NewProcessorFunc(
Specification{Name: "bench-metadata-field-set"},
func(_ context.Context, record opencdc.Record) (opencdc.Record, error) {
ref, err := resolver.Resolve(&record)
if err != nil {
return record, fmt.Errorf("failed to resolve reference: %w", err)
}
if err := ref.Set("bench-processor"); err != nil {
return record, fmt.Errorf("failed to set field: %w", err)
}
return record, nil
},
)
ctx := context.Background()
b.Run("Raw", func(b *testing.B) {
records := benchBatchRaw(benchBatchSize)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
out := proc.Process(ctx, records)
if len(out) != benchBatchSize {
b.Fatalf("expected %d processed records, got %d", benchBatchSize, len(out))
}
}
})
b.Run("Structured", func(b *testing.B) {
records := benchBatchStructured(benchBatchSize)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
out := proc.Process(ctx, records)
if len(out) != benchBatchSize {
b.Fatalf("expected %d processed records, got %d", benchBatchSize, len(out))
}
}
})
}
// BenchmarkProcessorFunc_Process_StructuredFieldSet benchmarks a processor
// that resolves a nested structured-payload reference and overwrites its
// value on every record in the batch, the pattern used by processors that
// rewrite a specific field deep in a record's payload (e.g. masking,
// normalization, enrichment).
func BenchmarkProcessorFunc_Process_StructuredFieldSet(b *testing.B) {
resolver, err := NewReferenceResolver(".Payload.After.customer.name")
if err != nil {
b.Fatal(err)
}
proc := NewProcessorFunc(
Specification{Name: "bench-structured-field-set"},
func(_ context.Context, record opencdc.Record) (opencdc.Record, error) {
ref, err := resolver.Resolve(&record)
if err != nil {
return record, fmt.Errorf("failed to resolve reference: %w", err)
}
if err := ref.Set("REDACTED"); err != nil {
return record, fmt.Errorf("failed to set field: %w", err)
}
return record, nil
},
)
ctx := context.Background()
records := benchBatchStructured(benchBatchSize)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
out := proc.Process(ctx, records)
if len(out) != benchBatchSize {
b.Fatalf("expected %d processed records, got %d", benchBatchSize, len(out))
}
}
}