Created
February 8, 2018 17:33
-
-
Save superfell/cff892d075b7585e46b5f1ec40206dcf to your computer and use it in GitHub Desktop.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
package main | |
import ( | |
"bytes" | |
"flag" | |
"fmt" | |
"log" | |
"os" | |
"sync" | |
"time" | |
sarama "gopkg.in/Shopify/sarama.v1" | |
) | |
const topic = "test2" | |
func main() { | |
broker := flag.String("b", "localhost:9092", "Kafka broker to connect to") | |
pcon := flag.Int("pc", 8, "Concurrency level of producer writes") | |
nummsg := flag.Int("msg", 50, "Number of test messages to create") | |
flag.Parse() | |
kconfig := sarama.NewConfig() | |
kconfig.Version = sarama.V0_11_0_0 | |
kconfig.Producer.RequiredAcks = sarama.WaitForAll | |
kconfig.Producer.Retry.Max = 10 | |
kconfig.Producer.Return.Successes = true | |
kconfig.Consumer.Return.Errors = true | |
kconfig.Consumer.Fetch.Min = 1 | |
kconfig.Consumer.Fetch.Default = 10 * 1024 * 1024 | |
kconfig.Consumer.Fetch.Max = 25 * 1024 * 1024 | |
client, err := sarama.NewClient([]string{*broker}, kconfig) | |
if err != nil { | |
log.Fatalf("Unable to create Kafka client: %v", err) | |
} | |
producer, err := sarama.NewSyncProducerFromClient(client) | |
if err != nil { | |
log.Fatalf("Unable to start kafka producer: %v", err) | |
} | |
consumer, err := sarama.NewConsumerFromClient(client) | |
if err != nil { | |
log.Fatalf("Unable to start kafka consumer: %v", err) | |
} | |
sarama.Logger = log.New(os.Stdout, "sarama:", 0) | |
offsetAtStart, err := client.GetOffset(topic, 0, sarama.OffsetNewest) | |
if err != nil { | |
log.Fatalf("Unable to read starting offset: %v", err) | |
} | |
workCh := make(chan []byte, *nummsg) | |
expectCh := make(chan expect, *nummsg) | |
go func() { | |
for i := 0; i < *nummsg; i++ { | |
data := []byte(fmt.Sprintf("%d", i)) | |
workCh <- data | |
} | |
close(workCh) | |
}() | |
pwg := sync.WaitGroup{} | |
for i := 0; i < *pcon; i++ { | |
pwg.Add(1) | |
go producerWorker(producer, workCh, expectCh, &pwg) | |
} | |
go func() { | |
pwg.Wait() | |
close(expectCh) | |
}() | |
writes := make(map[int64][]byte) | |
for e := range expectCh { | |
writes[e.offset] = e.value | |
} | |
producer.Close() | |
offsetAfterWrites, err := client.GetOffset(topic, 0, sarama.OffsetNewest) | |
fmt.Printf("Completed writing test messages, should be in offsets %d -%d\n", offsetAtStart, offsetAfterWrites) | |
// at this point all the test values have been produced, go comsume them and see if they match | |
pc, err := consumer.ConsumePartition(topic, 0, offsetAtStart) | |
if err != nil { | |
log.Fatalf("Unable to ConsumePartition: %v", err) | |
} | |
count := 0 | |
consumed := make(map[int64][]byte) | |
for { | |
select { | |
case e := <-pc.Errors(): | |
log.Fatalf("Error consuming topic: %v", e) | |
case m := <-pc.Messages(): | |
consumed[m.Offset] = m.Value | |
count++ | |
if count == *nummsg { | |
fmt.Printf("Consumed all %d messages we wrote\n", count) | |
dumpResults(writes, consumed, offsetAtStart, *nummsg) | |
return | |
} | |
case <-time.After(time.Second * 3): | |
fmt.Printf("Have read %d messages, giving up waiting for more\n", count) | |
dumpResults(writes, consumed, offsetAtStart, *nummsg) | |
return | |
} | |
} | |
} | |
func dumpResults(writes, consumed map[int64][]byte, first int64, num int) { | |
issues := 0 | |
last := first + int64(num) | |
for offset := first; offset < last; offset++ { | |
wrote := writes[offset] | |
cons, exists := consumed[offset] | |
if !exists { | |
fmt.Printf("Offset %d: wrote '%s', consumer never saw this offset\n", offset, wrote) | |
issues++ | |
} else { | |
good := bytes.Equal(wrote, cons) | |
fmt.Printf("Offset %d: wrote '%s', consumed '%s' good?:%v\n", offset, wrote, cons, good) | |
if !good { | |
issues++ | |
} | |
} | |
} | |
fmt.Printf("Processed %d messages, saw %d issues\n", num, issues) | |
} | |
func producerWorker(p sarama.SyncProducer, workCh chan []byte, resCh chan expect, wg *sync.WaitGroup) { | |
defer wg.Done() | |
for w := range workCh { | |
partition, offset, err := p.SendMessage(&sarama.ProducerMessage{ | |
Topic: topic, | |
Value: sarama.ByteEncoder(w), | |
}) | |
if err != nil { | |
log.Fatalf("Unable to write: %v", err) | |
} | |
if partition != 0 { | |
log.Fatalf("Expected partition to be 0, but was %d", partition) | |
} | |
resCh <- expect{ | |
offset: offset, | |
value: w, | |
} | |
} | |
} | |
type expect struct { | |
offset int64 | |
value []byte | |
} |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
results with a producer concurrency of 8, messages don't appear to be at the offset reported by the producer, the consumer see's holes in the offsets.