Documentation
¶
Index ¶
- Variables
- func NonPrefixPayloadExtractor(prefix string) func(wrapper MessageWrapper) map[string]interface{}
- func PrefixMetadataExtractor(prefix string) func(wrapper MessageWrapper) map[string]interface{}
- func PrefixPayloadExtractor(prefix string) func(wrapper MessageWrapper) map[string]interface{}
- func UnmarshallMapPayloadFromJson[T any](mapKey string, payloadType T) middleware.Middleware
- type IDGenerator
- type MessageUnmarshaller
- type MessageWrapper
- type Publisher
- type PublisherOption
- type PublisherOptions
- type StartPosition
- type Stream
- type Subscriber
- type SubscriberOption
- func WithConsumerGroup(group string) SubscriberOption
- func WithConsumerGroupCreateStreamIfMissing(create bool) SubscriberOption
- func WithConsumerGroupStartID(startID StartPosition) SubscriberOption
- func WithMetadataExtractor(extractor func(wrapper MessageWrapper) map[string]interface{}) SubscriberOption
- func WithPayloadExtractor(extractor func(wrapper MessageWrapper) map[string]interface{}) SubscriberOption
- func WithPendingMessageBatchSize(batchSize int) SubscriberOption
- func WithPendingMessageIdleTimeout(timeout time.Duration) SubscriberOption
- func WithProcessingTimeout(timeout time.Duration) SubscriberOption
- func WithProcessingTimeoutHandler(handler func(ctx context.Context, msg MessageWrapper)) SubscriberOption
- func WithStartID(startID StartPosition) SubscriberOption
- type SubscriberOptions
- type XAddValuesMarshaller
Constants ¶
This section is empty.
Variables ¶
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 MessageUnmarshaller ¶
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)
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 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
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 ¶
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.