amazon-kinesis-client-go/cmd/consumer/main.go

78 lines
2.1 KiB
Go
Raw Normal View History

2017-02-08 20:23:00 +00:00
package main
import (
"fmt"
"math/big"
"os"
"time"
"github.com/Clever/amazon-kinesis-client-go/kcl"
)
2017-07-18 19:19:40 +00:00
type sampleRecordProcessor struct {
checkpointer *kcl.Checkpointer
2017-02-08 20:23:00 +00:00
checkpointRetries int
checkpointFreq time.Duration
largestSeq *big.Int
largestSubSeq int
lastCheckpoint time.Time
}
2017-07-18 19:19:40 +00:00
func newSampleRecordProcessor() *sampleRecordProcessor {
return &sampleRecordProcessor{
2017-02-08 20:23:00 +00:00
checkpointRetries: 5,
checkpointFreq: 60 * time.Second,
}
}
2017-07-18 19:19:40 +00:00
func (srp *sampleRecordProcessor) Initialize(shardID string, checkpointer *kcl.Checkpointer) error {
2017-02-08 20:23:00 +00:00
srp.lastCheckpoint = time.Now()
srp.checkpointer = checkpointer
2017-02-08 20:23:00 +00:00
return nil
}
2017-07-18 19:19:40 +00:00
func (srp *sampleRecordProcessor) shouldUpdateSequence(sequenceNumber *big.Int, subSequenceNumber int) bool {
2017-02-08 20:23:00 +00:00
return srp.largestSeq == nil || sequenceNumber.Cmp(srp.largestSeq) == 1 ||
(sequenceNumber.Cmp(srp.largestSeq) == 0 && subSequenceNumber > srp.largestSubSeq)
}
2017-07-18 19:19:40 +00:00
func (srp *sampleRecordProcessor) ProcessRecords(records []kcl.Record) error {
2017-02-08 20:23:00 +00:00
for _, record := range records {
seqNumber := new(big.Int)
if _, ok := seqNumber.SetString(record.SequenceNumber, 10); !ok {
fmt.Fprintf(os.Stderr, "could not parse sequence number '%s'\n", record.SequenceNumber)
continue
}
if srp.shouldUpdateSequence(seqNumber, record.SubSequenceNumber) {
srp.largestSeq = seqNumber
srp.largestSubSeq = record.SubSequenceNumber
}
}
if time.Now().Sub(srp.lastCheckpoint) > srp.checkpointFreq {
largestSeq := srp.largestSeq.String()
srp.checkpointer.CheckpointWithRetry(&largestSeq, &srp.largestSubSeq, srp.checkpointRetries)
2017-02-08 20:23:00 +00:00
srp.lastCheckpoint = time.Now()
}
return nil
}
2017-07-18 19:19:40 +00:00
func (srp *sampleRecordProcessor) Shutdown(reason string) error {
2017-02-08 20:23:00 +00:00
if reason == "TERMINATE" {
fmt.Fprintf(os.Stderr, "Was told to terminate, will attempt to checkpoint.\n")
srp.checkpointer.Shutdown()
2017-02-08 20:23:00 +00:00
} else {
fmt.Fprintf(os.Stderr, "Shutting down due to failover. Will not checkpoint.\n")
}
return nil
}
func main() {
f, err := os.Create("/tmp/kcl_stderr")
if err != nil {
panic(err)
}
defer f.Close()
2017-07-18 19:19:40 +00:00
kclProcess := kcl.New(os.Stdin, os.Stdout, os.Stderr, newSampleRecordProcessor())
2017-02-08 20:23:00 +00:00
kclProcess.Run()
}