redis

package module
v0.0.0-...-7593c98 Latest Latest
Warning

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

Go to latest
Published: Aug 27, 2025 License: MIT Imports: 14 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var StreamDoesntExistErr = errors.New("stream does not exist")

Functions

func NonPrefixPayloadExtractor

func NonPrefixPayloadExtractor(prefix string) func(wrapper MessageWrapper) map[string]interface{}

NonPrefixPayloadExtractor is a function that extracts all values where the key doesn't starts with a prefix and use them as payload.

func PrefixMetadataExtractor

func PrefixMetadataExtractor(prefix string) func(wrapper MessageWrapper) map[string]interface{}

PrefixMetadataExtractor is a function that extracts all values where the key starts with a prefix and use them as metadata.

func PrefixPayloadExtractor

func PrefixPayloadExtractor(prefix string) func(wrapper MessageWrapper) map[string]interface{}

PrefixPayloadExtractor is a function that extracts all values where the key starts with a prefix and use them as payload.

func UnmarshallMapPayloadFromJson

func UnmarshallMapPayloadFromJson[T any](mapKey string, payloadType T) middleware.Middleware

Types

type IDGenerator

type IDGenerator func(msg *message.Message) (string, error)

type MessageUnmarshaller

type MessageUnmarshaller func(ctx context.Context, msg *redis.XMessage) (*message.Message, error)

type MessageWrapper

type MessageWrapper struct {
	// contains filtered or unexported fields
}

type Publisher

type Publisher struct {
	// contains filtered or unexported fields
}

func NewRedisPublisher

func NewRedisPublisher(
	c *redis.Client,
	routingFunc publisher.RoutingFunc,
	payloadMarshaller XAddValuesMarshaller,
	opts ...PublisherOption) (*Publisher, error)

func (*Publisher) Close

func (p *Publisher) Close() error

func (*Publisher) Publish

func (p *Publisher) Publish(message *message.Message) error

type PublisherOption

type PublisherOption func(*PublisherOptions)

func WithApprox

func WithApprox(approx bool) PublisherOption

WithApprox Approx is an optimization flag in redis stream, because it's costly to evict from a stream up to Maxlen. With approx=true, redis will evict up to Maxlen, but not exactly Maxlen, so there could be a bit more than Maxlen items in the stream. If not provided, approx=false is used.

func WithIDGenerator

func WithIDGenerator(idGen IDGenerator) PublisherOption

WithIDGenerator is a function that returns an ID for a given message to be send to redis. If not provided, redis will generate the ID.

func WithMaxlen

func WithMaxlen(maxlen int64) PublisherOption

WithMaxlen is a redis stream concept to limit the size of a stream. Old entries are evicted once maxlen is reached. If not provided, no limit is used.

func WithMetadataMarshaller

func WithMetadataMarshaller(marshaller XAddValuesMarshaller) PublisherOption

WithMetadataMarshaller is a function that returns a map of values representing the metadata part of the message to be send to redis. If not provided, the metadata string map of the message is used as is. This is useful if you want to process the metadata, for example, to add a prefix to all keys.

type PublisherOptions

type PublisherOptions struct {
	// contains filtered or unexported fields
}

PublisherOptions holds configuration options for a Publisher.

type StartPosition

type StartPosition string
const (
	StartFromBeginning StartPosition = "0"
	StartFromLatest    StartPosition = "$"
)

type Stream

type Stream struct {
	// contains filtered or unexported fields
}

func (Stream) Close

func (s Stream) Close() error

func (Stream) GetMessageID

func (s Stream) GetMessageID(message *redis.XAddArgs) string

func (Stream) Publish

func (s Stream) Publish(ctx context.Context, message *redis.XAddArgs) error

type Subscriber

type Subscriber struct {
	// contains filtered or unexported fields
}

func NewRedisSubscriber

func NewRedisSubscriber(
	c *redis.Client,
	stream string,
	opts ...SubscriberOption) (*Subscriber, error)

func (*Subscriber) Close

func (s *Subscriber) Close() error

func (*Subscriber) Subscribe

func (s *Subscriber) Subscribe(ctx context.Context) (<-chan *message.Message, error)

type SubscriberOption

type SubscriberOption func(*SubscriberOptions)

func WithConsumerGroup

func WithConsumerGroup(group string) SubscriberOption

WithConsumerGroup identifies the consumer group to which the subscriber belongs.

func WithConsumerGroupCreateStreamIfMissing

func WithConsumerGroupCreateStreamIfMissing(create bool) SubscriberOption

WithConsumerGroupCreateStreamIfMissing will create the stream if it does not exist in consumer group mode.

func WithConsumerGroupStartID

func WithConsumerGroupStartID(startID StartPosition) SubscriberOption

WithConsumerGroupStartID is the position the consumer of a consumer group should start from Using ConsumerGroupStartFromBeginning "0" means the consumer group will consume from the very first message. Using ConsumerGroupStartFromLatest "$" means the consumer group will consume from the latest message.

func WithMetadataExtractor

func WithMetadataExtractor(extractor func(wrapper MessageWrapper) map[string]interface{}) SubscriberOption

WithMetadataExtractor is a function that extracts metadata from a redis message. If not provided, no metadata will be extracted. The extractor usually needs to be aligned with the marshaller used by the publisher.

func WithPayloadExtractor

func WithPayloadExtractor(extractor func(wrapper MessageWrapper) map[string]interface{}) SubscriberOption

WithPayloadExtractor is a function that extracts payload from a redis message. If not provided, all the values in the redis message are converted into a map as payload.

func WithPendingMessageBatchSize

func WithPendingMessageBatchSize(batchSize int) SubscriberOption

WithPendingMessageBatchSize sets the maximum number of pending messages to check and claim per cycle. Default is 10. Only applies to consumer group mode.

func WithPendingMessageIdleTimeout

func WithPendingMessageIdleTimeout(timeout time.Duration) SubscriberOption

WithPendingMessageIdleTimeout sets how long a message must be idle before it can be claimed by another consumer. Default is 5 minutes. Only applies to consumer group mode.

func WithProcessingTimeout

func WithProcessingTimeout(timeout time.Duration) SubscriberOption

WithProcessingTimeout will dictate how long a message will be processed before it is nacked. 0 means no timeout, wait forever. Keep in mind that GCP uses the "Acknowledgement deadline" to determine if a message needs to be redelivered. ProcessingTimeout has no impact on the "Acknowledgement deadline". Default value is 600 seconds, which is the max value of the GCP "Acknowledgement deadline".

func WithProcessingTimeoutHandler

func WithProcessingTimeoutHandler(handler func(ctx context.Context, msg MessageWrapper)) SubscriberOption

WithProcessingTimeoutHandler takes a function that is called when a message processing times out.

func WithStartID

func WithStartID(startID StartPosition) SubscriberOption

WithStartID is the ID of the last message that was processed by the subscriber. This is used only when not using consumer groups. The default is "$" which means the subscriber will start from the latest message.

type SubscriberOptions

type SubscriberOptions struct {
	// contains filtered or unexported fields
}

type XAddValuesMarshaller

type XAddValuesMarshaller func(msg *message.Message) map[string]interface{}

func MarshallPayloadToJsonMap

func MarshallPayloadToJsonMap(mapKey string) XAddValuesMarshaller

func MetadataWithPrefix

func MetadataWithPrefix(prefix string) XAddValuesMarshaller

MetadataWithPrefix is a function that adds a prefix to all metadata keys.

Jump to

Keyboard shortcuts

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