apmqueue

package module
v2.10.4 Latest Latest
Warning

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

Go to latest
Published: Nov 6, 2024 License: Apache-2.0 Imports: 2 Imported by: 0

README

APM Queue

Producing and Consuming library that abstracts the details of producing and consuming records to and from Kafka.

Documentation

Overview

Package apmqueue provides an abstraction layer for producing and consuming Records from and to Kafka.

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrConsumerAlreadyRunning is returned by consumer.Run if it has already
	// been called.
	ErrConsumerAlreadyRunning = errors.New("consumer.Run: consumer already running")
)

Functions

This section is empty.

Types

type Consumer

type Consumer interface {
	// Run executes the consumer in a blocking manner. Returns
	// ErrConsumerAlreadyRunning when it has already been called.
	Run(ctx context.Context) error
	// Healthy returns an error if the consumer isn't healthy.
	Healthy(ctx context.Context) error
	// Close closes the consumer.
	Close() error
}

Consumer wraps the implementation details of the consumer implementation. Consumer implementations must support the defined delivery types.

type DeliveryType

type DeliveryType uint8

DeliveryType for the consumer. For more details See the supported DeliveryTypes.

const (
	// AtMostOnceDeliveryType acknowledges the message as soon as it's received
	// and decoded, without waiting for the message to be processed.
	AtMostOnceDeliveryType DeliveryType = iota
	// AtLeastOnceDeliveryType acknowledges the message after it has been
	// processed. It may or may not create duplicates, depending on how batches
	// are processed by the underlying Processor.
	AtLeastOnceDeliveryType
)

type Processor

type Processor interface {
	// Process processes one or more records within the passed context.
	// Process takes ownership of the passed records, callers must not mutate
	// a record after Process has been called.
	Process(context.Context, Record) error
}

Processor defines record processing signature.

type ProcessorFunc

type ProcessorFunc func(context.Context, Record) error

ProcessorFunc is a function type that implements the Processor interface.

func (ProcessorFunc) Process

func (f ProcessorFunc) Process(ctx context.Context, rs Record) error

Process returns f(ctx, records...).

type Producer

type Producer interface {
	// Produce produces N records. If the Producer is synchronous, waits until
	// all records are produced, otherwise, returns as soon as the records are
	// stored in the producer buffer, or when the records are produced to the
	// queue if sync producing is configured.
	// If the context has been enriched with metadata, each entry will be added
	// as a record's header.
	// Produce takes ownership of Record and any modifications after Produce is
	// called may cause an unhandled exception.
	Produce(ctx context.Context, rs ...Record) error
	// Healthy returns an error if the producer isn't healthy.
	Healthy(ctx context.Context) error
	// Close closes the producer.
	Close() error
}

Producer wraps the producer implementation details. Producer implementations must support sync and async production.

type Record

type Record struct {
	// OrderingKey is an optional field that is hashed to map to a partition.
	// Records with same ordering key are routed to the same partition.
	OrderingKey []byte
	// Value holds the record's content. It must not be mutated after Produce.
	Value []byte
	// Topics holds the topic where the record will be produced.
	Topic Topic
	// Partition identifies the partition ID where the record was polled from.
	// It is optional and only used for consumers.
	// When not specified, the zero value for int32 (0) identifies the only partition.
	Partition int32
}

Record wraps a record's value with the topic where it's produced / consumed.

type Topic

type Topic string

Topic represents a destination topic where to produce a message/record.

type TopicConsumer

type TopicConsumer struct {
	// Optional topic to monitor.
	Topic Topic
	// Optional regex expression to match topics for monitoring.
	Regex string
	// Required consumer name.
	Consumer string
}

TopicConsumer is used to monitor a set of consumer topics.

Directories

Path Synopsis
cmd
Package kafka abstracts the production and consumption of records to and from Kafka.
Package kafka abstracts the production and consumption of records to and from Kafka.
Package metrictest provides helpers for metric testing.
Package metrictest provides helpers for metric testing.
Package queuecontext provides convenient wrappers for storing and accessing a stored metadata.
Package queuecontext provides convenient wrappers for storing and accessing a stored metadata.
Package systemtest holds system/integration tests that test the apm-queue/kafka and apm-queue/pubsublite producer and consumer behaviour against real brokers/clusters.
Package systemtest holds system/integration tests that test the apm-queue/kafka and apm-queue/pubsublite producer and consumer behaviour against real brokers/clusters.

Jump to

Keyboard shortcuts

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