GitHub

File tree

  • component/async/kafka/group

  • test/docker/kafka

Original file line numberDiff line numberDiff line change

@@ -125,8 +125,8 @@ func (c *consumer) Consume(ctx context.Context) (<-chan async.Message, <-chan er

125125

closeConsumer(c.cg)

126126

return

127127

case consumerError := <-c.cg.Errors():

128-

closeConsumer(c.cg)

129128

chErr <- consumerError

129+

closeConsumer(c.cg)

130130

return

131131

}

132132

}

Original file line numberDiff line numberDiff line change

@@ -71,6 +71,7 @@ func TestGroupConsume_ClaimMessageError(t *testing.T) {

7171

chErr := make(chan error)

7272

go func() {

7373
74+

// Consumer will error out in ClaimMessage as no DecoderFunc has been set

7475

factory, err := group.New("test1", uuid.New().String(), []string{groupTopic2}, Brokers(),

7576

kafka.Version(sarama.V2_1_0_0.String()), kafka.StartFromNewest())

7677

if err != nil {

@@ -103,8 +104,10 @@ func TestGroupConsume_ClaimMessageError(t *testing.T) {

103104
104105

select {

105106

case <-chMessages:

106-

require.Fail(t, "no messages where expected")

107+

require.Fail(t, "no messages were expected")

107108

case err = <-chErr:

108-

require.EqualError(t, err, "kafka: tried to use a consumer group that was closed")

109+

require.EqualError(t, err, "kafka: error while consuming groupTopic2/0: "+

110+

"could not determine decoder failed to determine content type from message headers [] : "+

111+

"content type header is missing")

109112

}

110113

}

Read the original on github.com ↗