summaryrefslogtreecommitdiff
path: root/internal/parquet
diff options
context:
space:
mode:
authorPaul Buetow <paul@buetow.org>2026-03-13 07:46:29 +0200
committerPaul Buetow <paul@buetow.org>2026-03-13 07:46:29 +0200
commit973bc4be068d337ff9ab13c47d08485b1946d133 (patch)
tree8491c159fa052d632ec2d8866eae05b5669a36db /internal/parquet
parentb52fd1f297c178f17fe75f8fb03d5cbdd3ece71d (diff)
bench: add parquet recording benchmarks
Diffstat (limited to 'internal/parquet')
-rw-r--r--internal/parquet/bench_test.go68
1 files changed, 68 insertions, 0 deletions
diff --git a/internal/parquet/bench_test.go b/internal/parquet/bench_test.go
new file mode 100644
index 0000000..7083e27
--- /dev/null
+++ b/internal/parquet/bench_test.go
@@ -0,0 +1,68 @@
+package parquet
+
+import (
+ "path/filepath"
+ "testing"
+)
+
+var benchmarkRecordSink Record
+
+func BenchmarkWriterThroughput(b *testing.B) {
+ rows := benchmarkRecords(256)
+ dir := b.TempDir()
+ writer, err := NewWriter(filepath.Join(dir, "writer-throughput.parquet"), WriterConfig{}, FileMetadata{Mode: "bench"})
+ if err != nil {
+ b.Fatalf("NewWriter() error = %v", err)
+ }
+
+ b.ReportAllocs()
+ b.ResetTimer()
+
+ for i := 0; i < b.N; i++ {
+ if err := writer.WriteRows(rows); err != nil {
+ b.Fatalf("WriteRows() error = %v", err)
+ }
+ benchmarkRecordSink = rows[i%len(rows)]
+ }
+
+ b.StopTimer()
+ if err := writer.Close(); err != nil {
+ b.Fatalf("Close() error = %v", err)
+ }
+}
+
+func BenchmarkRecorderQueueHandoff(b *testing.B) {
+ row := testStreamRow(1, "read", false)
+ session := newRecordingSession(1)
+ recorder := &Recorder{
+ active: session,
+ status: Status{Active: true},
+ }
+
+ b.ReportAllocs()
+ b.ResetTimer()
+ for i := 0; i < b.N; i++ {
+ row.Seq = uint64(i + 1)
+ if err := recorder.Record(row, 0); err != nil {
+ b.Fatalf("Record() error = %v", err)
+ }
+ select {
+ case <-session.queue:
+ default:
+ b.Fatal("expected queued record request")
+ }
+ }
+
+ b.StopTimer()
+ session.stop(nil)
+ recorder.finishSession(session, nil)
+}
+
+func benchmarkRecords(n int) []Record {
+ rows := make([]Record, 0, n)
+ for i := 0; i < n; i++ {
+ row := testStreamRow(uint64(i+1), "read", i%7 == 0)
+ rows = append(rows, RecordFromStream(row, uint64(i%4)))
+ }
+ return rows
+}