|
@@ -300,7 +300,6 @@ func TestAsyncProducerFailureRetry(t *testing.T) {
|
|
|
for i := 0; i < 10; i++ {
|
|
for i := 0; i < 10; i++ {
|
|
|
producer.Input() <- &ProducerMessage{Topic: "my_topic", Key: nil, Value: StringEncoder(TestMessage)}
|
|
producer.Input() <- &ProducerMessage{Topic: "my_topic", Key: nil, Value: StringEncoder(TestMessage)}
|
|
|
}
|
|
}
|
|
|
- leader2.Returns(metadataLeader2)
|
|
|
|
|
leader2.Returns(prodSuccess)
|
|
leader2.Returns(prodSuccess)
|
|
|
expectResults(t, producer, 10, 0)
|
|
expectResults(t, producer, 10, 0)
|
|
|
|
|
|
|
@@ -468,7 +467,6 @@ func TestAsyncProducerMultipleRetries(t *testing.T) {
|
|
|
seedBroker.Returns(metadataLeader1)
|
|
seedBroker.Returns(metadataLeader1)
|
|
|
leader1.Returns(prodNotLeader)
|
|
leader1.Returns(prodNotLeader)
|
|
|
seedBroker.Returns(metadataLeader2)
|
|
seedBroker.Returns(metadataLeader2)
|
|
|
- seedBroker.Returns(metadataLeader2)
|
|
|
|
|
|
|
|
|
|
prodSuccess := new(ProduceResponse)
|
|
prodSuccess := new(ProduceResponse)
|
|
|
prodSuccess.AddTopicPartition("my_topic", 0, ErrNoError)
|
|
prodSuccess.AddTopicPartition("my_topic", 0, ErrNoError)
|
|
@@ -654,7 +652,6 @@ func TestAsyncProducerFlusherRetryCondition(t *testing.T) {
|
|
|
|
|
|
|
|
// succeed this time
|
|
// succeed this time
|
|
|
expectResults(t, producer, 5, 0)
|
|
expectResults(t, producer, 5, 0)
|
|
|
- seedBroker.Returns(metadataResponse)
|
|
|
|
|
|
|
|
|
|
// put five more through
|
|
// put five more through
|
|
|
for i := 0; i < 5; i++ {
|
|
for i := 0; i < 5; i++ {
|