pipelinelab

package
v1.0.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 7, 2026 License: Apache-2.0 Imports: 15 Imported by: 0

Documentation

Overview

Package pipelinelab is PipelineLab, an event pipeline: an in-process Kafka cluster (franz-go's kfake) with an orders topic, an enricher service that consumes it as the "enricher" consumer group, looks up each order's customer and writes the result to orders.enriched, and an admin API over HTTP that reports topics and consumer lag. It is the reference app for the event-pipelines pack.

Planted bottlenecks (see README.md), each switched off by a fix flag:

  • parallel: the enricher handles every record of every partition one at a time, so its throughput is one customer lookup at a time and a burst turns into lag (fix "parallel" works on each partition's records concurrently, keeping order within a partition).
  • commit: the enricher commits its offset after every record, a round trip to the group coordinator each (fix "commit" commits once per batch it polls).

Index

Constants

View Source
const (
	// Orders is the topic the pipeline consumes.
	Orders = "orders"
	// Enriched is the topic it writes to: same key, same partition, same
	// timestamp as the order, so a record's age there is the pipeline's
	// end-to-end delay.
	Enriched = "orders.enriched"
	// Group is the enricher's consumer group.
	Group = "enricher"
	// Partitions is each topic's partition count.
	Partitions = 6
	// Customers is how many customers the lookup knows: 1 to 10000.
	Customers = 10000
)

Variables

This section is empty.

Functions

func Open

func Open(cfg labkit.Config) (*labkit.App, error)

Open starts PipelineLab's Kafka cluster on cfg.Listen and its enricher, and returns the app.

Types

This section is empty.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL