Documentation ¶
Index ¶
- type LoadClientMessageProvider
- func (l *LoadClientMessageProvider) Close() error
- func (l *LoadClientMessageProvider) CommitOffsets(offsets map[int32]int64) error
- func (l *LoadClientMessageProvider) GetMessage(pollTimeout time.Duration) (*kafka.Message, error)
- func (l *LoadClientMessageProvider) SetRebalanceCallback(callback kafka.RebalanceCallback)
- func (l *LoadClientMessageProvider) Start() error
- func (l *LoadClientMessageProvider) Stop() error
- type LoadClientMessageProviderFactory
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type LoadClientMessageProvider ¶
type LoadClientMessageProvider struct {
// contains filtered or unexported fields
}
func (*LoadClientMessageProvider) Close ¶
func (l *LoadClientMessageProvider) Close() error
func (*LoadClientMessageProvider) CommitOffsets ¶
func (l *LoadClientMessageProvider) CommitOffsets(offsets map[int32]int64) error
func (*LoadClientMessageProvider) GetMessage ¶
func (*LoadClientMessageProvider) SetRebalanceCallback ¶
func (l *LoadClientMessageProvider) SetRebalanceCallback(callback kafka.RebalanceCallback)
func (*LoadClientMessageProvider) Start ¶
func (l *LoadClientMessageProvider) Start() error
func (*LoadClientMessageProvider) Stop ¶
func (l *LoadClientMessageProvider) Stop() error
type LoadClientMessageProviderFactory ¶
type LoadClientMessageProviderFactory struct {
// contains filtered or unexported fields
}
func (*LoadClientMessageProviderFactory) NewMessageProvider ¶
func (l *LoadClientMessageProviderFactory) NewMessageProvider() (kafka.MessageProvider, error)
Click to show internal directories.
Click to hide internal directories.