Documentation ¶
Index ¶
- func WaitManyAckQueues(queues ...*AckQueue) bool
- type AckQueue
- func (q *AckQueue) Dispose()
- func (q *AckQueue) Disposed() bool
- func (q *AckQueue) Get() (utils.MyULID, conf.DestinationType, error)
- func (q *AckQueue) GetMany(max uint32) (res []UidDest)
- func (q *AckQueue) GetManyInto(uids *[]UidDest)
- func (q *AckQueue) Has() bool
- func (q *AckQueue) Put(uid utils.MyULID, dest conf.DestinationType) error
- func (q *AckQueue) Wait() bool
- type BSliceQueue
- func (q *BSliceQueue) Dispose()
- func (q *BSliceQueue) Disposed() bool
- func (q *BSliceQueue) Get() (bool, utils.MyULID, string, error)
- func (q *BSliceQueue) GetMany(max uint32) (slices []string)
- func (q *BSliceQueue) GetManyInto(slices *[]string, uids *[]utils.MyULID)
- func (q *BSliceQueue) GetManyIntoMap(m *map[utils.MyULID]string, max uint32)
- func (q *BSliceQueue) GetManySlicesInto(slices *[]string)
- func (q *BSliceQueue) Has() bool
- func (q *BSliceQueue) Put(uid utils.MyULID, m string) error
- func (q *BSliceQueue) PutSlice(m string) error
- func (q *BSliceQueue) Wait(timeout time.Duration) bool
- type Data
- type KafkaProducerAck
- type KafkaProducerAckQueue
- func (q *KafkaProducerAckQueue) Dispose()
- func (q *KafkaProducerAckQueue) Disposed() bool
- func (q *KafkaProducerAckQueue) Get() (KafkaProducerAck, error)
- func (q *KafkaProducerAckQueue) Has() bool
- func (q *KafkaProducerAckQueue) Put(offset int64, partition int32, topic string) error
- func (q *KafkaProducerAckQueue) Wait() bool
- type KafkaQueues
- type Ring
- func (rb *Ring) Cap() uint64
- func (rb *Ring) Dispose()
- func (rb *Ring) Get() (Data, error)
- func (rb *Ring) IsDisposed() bool
- func (rb *Ring) Len() uint64
- func (rb *Ring) Offer(item Data) (bool, error)
- func (rb *Ring) Poll(timeout time.Duration) (Data, error)
- func (rb *Ring) PollDeadline(deadline time.Time) (Data, error)
- func (rb *Ring) Put(item Data) error
- type TopicPartition
- type UidDest
- type WrappedQueue
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func WaitManyAckQueues ¶
Types ¶
type AckQueue ¶
type AckQueue struct {
// contains filtered or unexported fields
}
func NewAckQueue ¶
func NewAckQueue() *AckQueue
func (*AckQueue) GetManyInto ¶
type BSliceQueue ¶
type BSliceQueue struct {
// contains filtered or unexported fields
}
func NewBSliceQueue ¶
func NewBSliceQueue() *BSliceQueue
func (*BSliceQueue) Dispose ¶
func (q *BSliceQueue) Dispose()
func (*BSliceQueue) Disposed ¶
func (q *BSliceQueue) Disposed() bool
func (*BSliceQueue) GetMany ¶
func (q *BSliceQueue) GetMany(max uint32) (slices []string)
func (*BSliceQueue) GetManyInto ¶
func (q *BSliceQueue) GetManyInto(slices *[]string, uids *[]utils.MyULID)
func (*BSliceQueue) GetManyIntoMap ¶
func (q *BSliceQueue) GetManyIntoMap(m *map[utils.MyULID]string, max uint32)
func (*BSliceQueue) GetManySlicesInto ¶
func (q *BSliceQueue) GetManySlicesInto(slices *[]string)
func (*BSliceQueue) Has ¶
func (q *BSliceQueue) Has() bool
func (*BSliceQueue) PutSlice ¶
func (q *BSliceQueue) PutSlice(m string) error
type KafkaProducerAck ¶
type KafkaProducerAck struct { Offset int64 TopicPartition }
type KafkaProducerAckQueue ¶
type KafkaProducerAckQueue struct {
// contains filtered or unexported fields
}
func NewKafkaProducerAckQueue ¶
func NewKafkaProducerAckQueue() *KafkaProducerAckQueue
func (*KafkaProducerAckQueue) Dispose ¶
func (q *KafkaProducerAckQueue) Dispose()
func (*KafkaProducerAckQueue) Disposed ¶
func (q *KafkaProducerAckQueue) Disposed() bool
func (*KafkaProducerAckQueue) Get ¶
func (q *KafkaProducerAckQueue) Get() (KafkaProducerAck, error)
func (*KafkaProducerAckQueue) Has ¶
func (q *KafkaProducerAckQueue) Has() bool
func (*KafkaProducerAckQueue) Put ¶
func (q *KafkaProducerAckQueue) Put(offset int64, partition int32, topic string) error
func (*KafkaProducerAckQueue) Wait ¶
func (q *KafkaProducerAckQueue) Wait() bool
type KafkaQueues ¶
type KafkaQueues struct {
// contains filtered or unexported fields
}
func NewQueueFactory ¶
func NewQueueFactory() *KafkaQueues
func (*KafkaQueues) Delete ¶
func (qs *KafkaQueues) Delete(q *WrappedQueue)
func (*KafkaQueues) Get ¶
func (qs *KafkaQueues) Get(qid uint32) *WrappedQueue
func (*KafkaQueues) New ¶
func (qs *KafkaQueues) New() (q *WrappedQueue)
type Ring ¶
type Ring struct {
// contains filtered or unexported fields
}
Ring is a thread-safe bounded queue that stores Data messages.
func (*Ring) Dispose ¶
func (rb *Ring) Dispose()
Dispose will dispose of this queue and free any blocked threads in the Put and/or Get methods. Calling those methods on a disposed queue will return an error.
func (*Ring) Get ¶
Get will return the next item in the queue. This call will block if the queue is empty. This call will unblock when an item is added to the queue or Dispose is called on the queue. An error will be returned if the queue is disposed.
func (*Ring) IsDisposed ¶
IsDisposed will return a bool indicating if this queue has been disposed.
func (*Ring) Offer ¶
Offer adds the provided item to the queue if there is space. If the queue is full, this call will return false. An error will be returned if the queue is disposed.
func (*Ring) Poll ¶
Poll will return the next item in the queue. This call will block if the queue is empty. This call will unblock when an item is added to the queue, Dispose is called on the queue, or the timeout is reached. An error will be returned if the queue is disposed or a timeout occurs. A non-positive timeout will block indefinitely.
type TopicPartition ¶
type WrappedQueue ¶
type WrappedQueue struct { *KafkaProducerAckQueue // contains filtered or unexported fields }
func (*WrappedQueue) ID ¶
func (q *WrappedQueue) ID() uint32