Versions in this module Expand all Collapse all v0 v0.1.0 Oct 16, 2017 Changes in this version + type Buffer struct + MaxRecordCount int + func (b *Buffer) AddRecord(r *kinesis.Record) + func (b *Buffer) FirstSeq() string + func (b *Buffer) Flush() + func (b *Buffer) GetRecords() []*kinesis.Record + func (b *Buffer) LastSeq() string + func (b *Buffer) RecordCount() int + func (b *Buffer) ShardID() string + func (b *Buffer) ShouldFlush() bool + type Checkpoint interface + CheckpointExists func(shardID string) bool + SequenceNumber func() string + SetCheckpoint func(shardID string, sequenceNumber string) + type Config struct + AppName string + BufferSize int + Checkpoint Checkpoint + FlushInterval time.Duration + Logger log.Interface + StreamName string + type Consumer struct + func NewConsumer(config Config) *Consumer + func (c *Consumer) Start(handler Handler) + type Handler interface + HandleRecords func(b Buffer) + type HandlerFunc func(b Buffer) + func (h HandlerFunc) HandleRecords(b Buffer) + type RedisCheckpoint struct + AppName string + StreamName string + func (c *RedisCheckpoint) CheckpointExists(shardID string) bool + func (c *RedisCheckpoint) SequenceNumber() string + func (c *RedisCheckpoint) SetCheckpoint(shardID string, sequenceNumber string)