consumer_test.go 39 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344
  1. package sarama
  2. import (
  3. "log"
  4. "os"
  5. "os/signal"
  6. "reflect"
  7. "sync"
  8. "sync/atomic"
  9. "testing"
  10. "time"
  11. )
  12. var testMsg = StringEncoder("Foo")
  13. // If a particular offset is provided then messages are consumed starting from
  14. // that offset.
  15. func TestConsumerOffsetManual(t *testing.T) {
  16. // Given
  17. broker0 := NewMockBroker(t, 0)
  18. mockFetchResponse := NewMockFetchResponse(t, 1)
  19. for i := 0; i < 10; i++ {
  20. mockFetchResponse.SetMessage("my_topic", 0, int64(i+1234), testMsg)
  21. }
  22. broker0.SetHandlerByMap(map[string]MockResponse{
  23. "MetadataRequest": NewMockMetadataResponse(t).
  24. SetBroker(broker0.Addr(), broker0.BrokerID()).
  25. SetLeader("my_topic", 0, broker0.BrokerID()),
  26. "OffsetRequest": NewMockOffsetResponse(t).
  27. SetOffset("my_topic", 0, OffsetOldest, 0).
  28. SetOffset("my_topic", 0, OffsetNewest, 2345),
  29. "FetchRequest": mockFetchResponse,
  30. })
  31. // When
  32. master, err := NewConsumer([]string{broker0.Addr()}, nil)
  33. if err != nil {
  34. t.Fatal(err)
  35. }
  36. consumer, err := master.ConsumePartition("my_topic", 0, 1234)
  37. if err != nil {
  38. t.Fatal(err)
  39. }
  40. // Then: messages starting from offset 1234 are consumed.
  41. for i := 0; i < 10; i++ {
  42. select {
  43. case message := <-consumer.Messages():
  44. assertMessageOffset(t, message, int64(i+1234))
  45. case err := <-consumer.Errors():
  46. t.Error(err)
  47. }
  48. }
  49. safeClose(t, consumer)
  50. safeClose(t, master)
  51. broker0.Close()
  52. }
  53. // If `OffsetNewest` is passed as the initial offset then the first consumed
  54. // message is indeed corresponds to the offset that broker claims to be the
  55. // newest in its metadata response.
  56. func TestConsumerOffsetNewest(t *testing.T) {
  57. // Given
  58. broker0 := NewMockBroker(t, 0)
  59. broker0.SetHandlerByMap(map[string]MockResponse{
  60. "MetadataRequest": NewMockMetadataResponse(t).
  61. SetBroker(broker0.Addr(), broker0.BrokerID()).
  62. SetLeader("my_topic", 0, broker0.BrokerID()),
  63. "OffsetRequest": NewMockOffsetResponse(t).
  64. SetOffset("my_topic", 0, OffsetNewest, 10).
  65. SetOffset("my_topic", 0, OffsetOldest, 7),
  66. "FetchRequest": NewMockFetchResponse(t, 1).
  67. SetMessage("my_topic", 0, 9, testMsg).
  68. SetMessage("my_topic", 0, 10, testMsg).
  69. SetMessage("my_topic", 0, 11, testMsg).
  70. SetHighWaterMark("my_topic", 0, 14),
  71. })
  72. master, err := NewConsumer([]string{broker0.Addr()}, nil)
  73. if err != nil {
  74. t.Fatal(err)
  75. }
  76. // When
  77. consumer, err := master.ConsumePartition("my_topic", 0, OffsetNewest)
  78. if err != nil {
  79. t.Fatal(err)
  80. }
  81. // Then
  82. assertMessageOffset(t, <-consumer.Messages(), 10)
  83. if hwmo := consumer.HighWaterMarkOffset(); hwmo != 14 {
  84. t.Errorf("Expected high water mark offset 14, found %d", hwmo)
  85. }
  86. safeClose(t, consumer)
  87. safeClose(t, master)
  88. broker0.Close()
  89. }
  90. // It is possible to close a partition consumer and create the same anew.
  91. func TestConsumerRecreate(t *testing.T) {
  92. // Given
  93. broker0 := NewMockBroker(t, 0)
  94. broker0.SetHandlerByMap(map[string]MockResponse{
  95. "MetadataRequest": NewMockMetadataResponse(t).
  96. SetBroker(broker0.Addr(), broker0.BrokerID()).
  97. SetLeader("my_topic", 0, broker0.BrokerID()),
  98. "OffsetRequest": NewMockOffsetResponse(t).
  99. SetOffset("my_topic", 0, OffsetOldest, 0).
  100. SetOffset("my_topic", 0, OffsetNewest, 1000),
  101. "FetchRequest": NewMockFetchResponse(t, 1).
  102. SetMessage("my_topic", 0, 10, testMsg),
  103. })
  104. c, err := NewConsumer([]string{broker0.Addr()}, nil)
  105. if err != nil {
  106. t.Fatal(err)
  107. }
  108. pc, err := c.ConsumePartition("my_topic", 0, 10)
  109. if err != nil {
  110. t.Fatal(err)
  111. }
  112. assertMessageOffset(t, <-pc.Messages(), 10)
  113. // When
  114. safeClose(t, pc)
  115. pc, err = c.ConsumePartition("my_topic", 0, 10)
  116. if err != nil {
  117. t.Fatal(err)
  118. }
  119. // Then
  120. assertMessageOffset(t, <-pc.Messages(), 10)
  121. safeClose(t, pc)
  122. safeClose(t, c)
  123. broker0.Close()
  124. }
  125. // An attempt to consume the same partition twice should fail.
  126. func TestConsumerDuplicate(t *testing.T) {
  127. // Given
  128. broker0 := NewMockBroker(t, 0)
  129. broker0.SetHandlerByMap(map[string]MockResponse{
  130. "MetadataRequest": NewMockMetadataResponse(t).
  131. SetBroker(broker0.Addr(), broker0.BrokerID()).
  132. SetLeader("my_topic", 0, broker0.BrokerID()),
  133. "OffsetRequest": NewMockOffsetResponse(t).
  134. SetOffset("my_topic", 0, OffsetOldest, 0).
  135. SetOffset("my_topic", 0, OffsetNewest, 1000),
  136. "FetchRequest": NewMockFetchResponse(t, 1),
  137. })
  138. config := NewConfig()
  139. config.ChannelBufferSize = 0
  140. c, err := NewConsumer([]string{broker0.Addr()}, config)
  141. if err != nil {
  142. t.Fatal(err)
  143. }
  144. pc1, err := c.ConsumePartition("my_topic", 0, 0)
  145. if err != nil {
  146. t.Fatal(err)
  147. }
  148. // When
  149. pc2, err := c.ConsumePartition("my_topic", 0, 0)
  150. // Then
  151. if pc2 != nil || err != ConfigurationError("That topic/partition is already being consumed") {
  152. t.Fatal("A partition cannot be consumed twice at the same time")
  153. }
  154. safeClose(t, pc1)
  155. safeClose(t, c)
  156. broker0.Close()
  157. }
  158. func runConsumerLeaderRefreshErrorTestWithConfig(t *testing.T, config *Config) {
  159. // Given
  160. broker0 := NewMockBroker(t, 100)
  161. // Stage 1: my_topic/0 served by broker0
  162. Logger.Printf(" STAGE 1")
  163. broker0.SetHandlerByMap(map[string]MockResponse{
  164. "MetadataRequest": NewMockMetadataResponse(t).
  165. SetBroker(broker0.Addr(), broker0.BrokerID()).
  166. SetLeader("my_topic", 0, broker0.BrokerID()),
  167. "OffsetRequest": NewMockOffsetResponse(t).
  168. SetOffset("my_topic", 0, OffsetOldest, 123).
  169. SetOffset("my_topic", 0, OffsetNewest, 1000),
  170. "FetchRequest": NewMockFetchResponse(t, 1).
  171. SetMessage("my_topic", 0, 123, testMsg),
  172. })
  173. c, err := NewConsumer([]string{broker0.Addr()}, config)
  174. if err != nil {
  175. t.Fatal(err)
  176. }
  177. pc, err := c.ConsumePartition("my_topic", 0, OffsetOldest)
  178. if err != nil {
  179. t.Fatal(err)
  180. }
  181. assertMessageOffset(t, <-pc.Messages(), 123)
  182. // Stage 2: broker0 says that it is no longer the leader for my_topic/0,
  183. // but the requests to retrieve metadata fail with network timeout.
  184. Logger.Printf(" STAGE 2")
  185. fetchResponse2 := &FetchResponse{}
  186. fetchResponse2.AddError("my_topic", 0, ErrNotLeaderForPartition)
  187. broker0.SetHandlerByMap(map[string]MockResponse{
  188. "FetchRequest": NewMockWrapper(fetchResponse2),
  189. })
  190. if consErr := <-pc.Errors(); consErr.Err != ErrOutOfBrokers {
  191. t.Errorf("Unexpected error: %v", consErr.Err)
  192. }
  193. // Stage 3: finally the metadata returned by broker0 tells that broker1 is
  194. // a new leader for my_topic/0. Consumption resumes.
  195. Logger.Printf(" STAGE 3")
  196. broker1 := NewMockBroker(t, 101)
  197. broker1.SetHandlerByMap(map[string]MockResponse{
  198. "FetchRequest": NewMockFetchResponse(t, 1).
  199. SetMessage("my_topic", 0, 124, testMsg),
  200. })
  201. broker0.SetHandlerByMap(map[string]MockResponse{
  202. "MetadataRequest": NewMockMetadataResponse(t).
  203. SetBroker(broker0.Addr(), broker0.BrokerID()).
  204. SetBroker(broker1.Addr(), broker1.BrokerID()).
  205. SetLeader("my_topic", 0, broker1.BrokerID()),
  206. })
  207. assertMessageOffset(t, <-pc.Messages(), 124)
  208. safeClose(t, pc)
  209. safeClose(t, c)
  210. broker1.Close()
  211. broker0.Close()
  212. }
  213. // If consumer fails to refresh metadata it keeps retrying with frequency
  214. // specified by `Config.Consumer.Retry.Backoff`.
  215. func TestConsumerLeaderRefreshError(t *testing.T) {
  216. config := NewConfig()
  217. config.Net.ReadTimeout = 100 * time.Millisecond
  218. config.Consumer.Retry.Backoff = 200 * time.Millisecond
  219. config.Consumer.Return.Errors = true
  220. config.Metadata.Retry.Max = 0
  221. runConsumerLeaderRefreshErrorTestWithConfig(t, config)
  222. }
  223. func TestConsumerLeaderRefreshErrorWithBackoffFunc(t *testing.T) {
  224. var calls int32 = 0
  225. config := NewConfig()
  226. config.Net.ReadTimeout = 100 * time.Millisecond
  227. config.Consumer.Retry.BackoffFunc = func(retries int) time.Duration {
  228. atomic.AddInt32(&calls, 1)
  229. return 200 * time.Millisecond
  230. }
  231. config.Consumer.Return.Errors = true
  232. config.Metadata.Retry.Max = 0
  233. runConsumerLeaderRefreshErrorTestWithConfig(t, config)
  234. // we expect at least one call to our backoff function
  235. if calls == 0 {
  236. t.Fail()
  237. }
  238. }
  239. func TestConsumerInvalidTopic(t *testing.T) {
  240. // Given
  241. broker0 := NewMockBroker(t, 100)
  242. broker0.SetHandlerByMap(map[string]MockResponse{
  243. "MetadataRequest": NewMockMetadataResponse(t).
  244. SetBroker(broker0.Addr(), broker0.BrokerID()),
  245. })
  246. c, err := NewConsumer([]string{broker0.Addr()}, nil)
  247. if err != nil {
  248. t.Fatal(err)
  249. }
  250. // When
  251. pc, err := c.ConsumePartition("my_topic", 0, OffsetOldest)
  252. // Then
  253. if pc != nil || err != ErrUnknownTopicOrPartition {
  254. t.Errorf("Should fail with, err=%v", err)
  255. }
  256. safeClose(t, c)
  257. broker0.Close()
  258. }
  259. // Nothing bad happens if a partition consumer that has no leader assigned at
  260. // the moment is closed.
  261. func TestConsumerClosePartitionWithoutLeader(t *testing.T) {
  262. // Given
  263. broker0 := NewMockBroker(t, 100)
  264. broker0.SetHandlerByMap(map[string]MockResponse{
  265. "MetadataRequest": NewMockMetadataResponse(t).
  266. SetBroker(broker0.Addr(), broker0.BrokerID()).
  267. SetLeader("my_topic", 0, broker0.BrokerID()),
  268. "OffsetRequest": NewMockOffsetResponse(t).
  269. SetOffset("my_topic", 0, OffsetOldest, 123).
  270. SetOffset("my_topic", 0, OffsetNewest, 1000),
  271. "FetchRequest": NewMockFetchResponse(t, 1).
  272. SetMessage("my_topic", 0, 123, testMsg),
  273. })
  274. config := NewConfig()
  275. config.Net.ReadTimeout = 100 * time.Millisecond
  276. config.Consumer.Retry.Backoff = 100 * time.Millisecond
  277. config.Consumer.Return.Errors = true
  278. config.Metadata.Retry.Max = 0
  279. c, err := NewConsumer([]string{broker0.Addr()}, config)
  280. if err != nil {
  281. t.Fatal(err)
  282. }
  283. pc, err := c.ConsumePartition("my_topic", 0, OffsetOldest)
  284. if err != nil {
  285. t.Fatal(err)
  286. }
  287. assertMessageOffset(t, <-pc.Messages(), 123)
  288. // broker0 says that it is no longer the leader for my_topic/0, but the
  289. // requests to retrieve metadata fail with network timeout.
  290. fetchResponse2 := &FetchResponse{}
  291. fetchResponse2.AddError("my_topic", 0, ErrNotLeaderForPartition)
  292. broker0.SetHandlerByMap(map[string]MockResponse{
  293. "FetchRequest": NewMockWrapper(fetchResponse2),
  294. })
  295. // When
  296. if consErr := <-pc.Errors(); consErr.Err != ErrOutOfBrokers {
  297. t.Errorf("Unexpected error: %v", consErr.Err)
  298. }
  299. // Then: the partition consumer can be closed without any problem.
  300. safeClose(t, pc)
  301. safeClose(t, c)
  302. broker0.Close()
  303. }
  304. // If the initial offset passed on partition consumer creation is out of the
  305. // actual offset range for the partition, then the partition consumer stops
  306. // immediately closing its output channels.
  307. func TestConsumerShutsDownOutOfRange(t *testing.T) {
  308. // Given
  309. broker0 := NewMockBroker(t, 0)
  310. fetchResponse := new(FetchResponse)
  311. fetchResponse.AddError("my_topic", 0, ErrOffsetOutOfRange)
  312. broker0.SetHandlerByMap(map[string]MockResponse{
  313. "MetadataRequest": NewMockMetadataResponse(t).
  314. SetBroker(broker0.Addr(), broker0.BrokerID()).
  315. SetLeader("my_topic", 0, broker0.BrokerID()),
  316. "OffsetRequest": NewMockOffsetResponse(t).
  317. SetOffset("my_topic", 0, OffsetNewest, 1234).
  318. SetOffset("my_topic", 0, OffsetOldest, 7),
  319. "FetchRequest": NewMockWrapper(fetchResponse),
  320. })
  321. master, err := NewConsumer([]string{broker0.Addr()}, nil)
  322. if err != nil {
  323. t.Fatal(err)
  324. }
  325. // When
  326. consumer, err := master.ConsumePartition("my_topic", 0, 101)
  327. if err != nil {
  328. t.Fatal(err)
  329. }
  330. // Then: consumer should shut down closing its messages and errors channels.
  331. if _, ok := <-consumer.Messages(); ok {
  332. t.Error("Expected the consumer to shut down")
  333. }
  334. safeClose(t, consumer)
  335. safeClose(t, master)
  336. broker0.Close()
  337. }
  338. // If a fetch response contains messages with offsets that are smaller then
  339. // requested, then such messages are ignored.
  340. func TestConsumerExtraOffsets(t *testing.T) {
  341. // Given
  342. legacyFetchResponse := &FetchResponse{}
  343. legacyFetchResponse.AddMessage("my_topic", 0, nil, testMsg, 1)
  344. legacyFetchResponse.AddMessage("my_topic", 0, nil, testMsg, 2)
  345. legacyFetchResponse.AddMessage("my_topic", 0, nil, testMsg, 3)
  346. legacyFetchResponse.AddMessage("my_topic", 0, nil, testMsg, 4)
  347. newFetchResponse := &FetchResponse{Version: 4}
  348. newFetchResponse.AddRecord("my_topic", 0, nil, testMsg, 1)
  349. newFetchResponse.AddRecord("my_topic", 0, nil, testMsg, 2)
  350. newFetchResponse.AddRecord("my_topic", 0, nil, testMsg, 3)
  351. newFetchResponse.AddRecord("my_topic", 0, nil, testMsg, 4)
  352. newFetchResponse.SetLastOffsetDelta("my_topic", 0, 4)
  353. newFetchResponse.SetLastStableOffset("my_topic", 0, 4)
  354. for _, fetchResponse1 := range []*FetchResponse{legacyFetchResponse, newFetchResponse} {
  355. var offsetResponseVersion int16
  356. cfg := NewConfig()
  357. cfg.Consumer.Return.Errors = true
  358. if fetchResponse1.Version >= 4 {
  359. cfg.Version = V0_11_0_0
  360. offsetResponseVersion = 1
  361. }
  362. broker0 := NewMockBroker(t, 0)
  363. fetchResponse2 := &FetchResponse{}
  364. fetchResponse2.Version = fetchResponse1.Version
  365. fetchResponse2.AddError("my_topic", 0, ErrNoError)
  366. broker0.SetHandlerByMap(map[string]MockResponse{
  367. "MetadataRequest": NewMockMetadataResponse(t).
  368. SetBroker(broker0.Addr(), broker0.BrokerID()).
  369. SetLeader("my_topic", 0, broker0.BrokerID()),
  370. "OffsetRequest": NewMockOffsetResponse(t).
  371. SetVersion(offsetResponseVersion).
  372. SetOffset("my_topic", 0, OffsetNewest, 1234).
  373. SetOffset("my_topic", 0, OffsetOldest, 0),
  374. "FetchRequest": NewMockSequence(fetchResponse1, fetchResponse2),
  375. })
  376. master, err := NewConsumer([]string{broker0.Addr()}, cfg)
  377. if err != nil {
  378. t.Fatal(err)
  379. }
  380. // When
  381. consumer, err := master.ConsumePartition("my_topic", 0, 3)
  382. if err != nil {
  383. t.Fatal(err)
  384. }
  385. // Then: messages with offsets 1 and 2 are not returned even though they
  386. // are present in the response.
  387. select {
  388. case msg := <-consumer.Messages():
  389. assertMessageOffset(t, msg, 3)
  390. case err := <-consumer.Errors():
  391. t.Fatal(err)
  392. }
  393. select {
  394. case msg := <-consumer.Messages():
  395. assertMessageOffset(t, msg, 4)
  396. case err := <-consumer.Errors():
  397. t.Fatal(err)
  398. }
  399. safeClose(t, consumer)
  400. safeClose(t, master)
  401. broker0.Close()
  402. }
  403. }
  404. // In some situations broker may return a block containing only
  405. // messages older then requested, even though there would be
  406. // more messages if higher offset was requested.
  407. func TestConsumerReceivingFetchResponseWithTooOldRecords(t *testing.T) {
  408. // Given
  409. fetchResponse1 := &FetchResponse{Version: 4}
  410. fetchResponse1.AddRecord("my_topic", 0, nil, testMsg, 1)
  411. fetchResponse2 := &FetchResponse{Version: 4}
  412. fetchResponse2.AddRecord("my_topic", 0, nil, testMsg, 1000000)
  413. cfg := NewConfig()
  414. cfg.Consumer.Return.Errors = true
  415. cfg.Version = V0_11_0_0
  416. broker0 := NewMockBroker(t, 0)
  417. broker0.SetHandlerByMap(map[string]MockResponse{
  418. "MetadataRequest": NewMockMetadataResponse(t).
  419. SetBroker(broker0.Addr(), broker0.BrokerID()).
  420. SetLeader("my_topic", 0, broker0.BrokerID()),
  421. "OffsetRequest": NewMockOffsetResponse(t).
  422. SetVersion(1).
  423. SetOffset("my_topic", 0, OffsetNewest, 1234).
  424. SetOffset("my_topic", 0, OffsetOldest, 0),
  425. "FetchRequest": NewMockSequence(fetchResponse1, fetchResponse2),
  426. })
  427. master, err := NewConsumer([]string{broker0.Addr()}, cfg)
  428. if err != nil {
  429. t.Fatal(err)
  430. }
  431. // When
  432. consumer, err := master.ConsumePartition("my_topic", 0, 2)
  433. if err != nil {
  434. t.Fatal(err)
  435. }
  436. select {
  437. case msg := <-consumer.Messages():
  438. assertMessageOffset(t, msg, 1000000)
  439. case err := <-consumer.Errors():
  440. t.Fatal(err)
  441. }
  442. safeClose(t, consumer)
  443. safeClose(t, master)
  444. broker0.Close()
  445. }
  446. func TestConsumeMessageWithNewerFetchAPIVersion(t *testing.T) {
  447. // Given
  448. fetchResponse1 := &FetchResponse{Version: 4}
  449. fetchResponse1.AddMessage("my_topic", 0, nil, testMsg, 1)
  450. fetchResponse1.AddMessage("my_topic", 0, nil, testMsg, 2)
  451. cfg := NewConfig()
  452. cfg.Version = V0_11_0_0
  453. broker0 := NewMockBroker(t, 0)
  454. fetchResponse2 := &FetchResponse{}
  455. fetchResponse2.Version = 4
  456. fetchResponse2.AddError("my_topic", 0, ErrNoError)
  457. broker0.SetHandlerByMap(map[string]MockResponse{
  458. "MetadataRequest": NewMockMetadataResponse(t).
  459. SetBroker(broker0.Addr(), broker0.BrokerID()).
  460. SetLeader("my_topic", 0, broker0.BrokerID()),
  461. "OffsetRequest": NewMockOffsetResponse(t).
  462. SetVersion(1).
  463. SetOffset("my_topic", 0, OffsetNewest, 1234).
  464. SetOffset("my_topic", 0, OffsetOldest, 0),
  465. "FetchRequest": NewMockSequence(fetchResponse1, fetchResponse2),
  466. })
  467. master, err := NewConsumer([]string{broker0.Addr()}, cfg)
  468. if err != nil {
  469. t.Fatal(err)
  470. }
  471. // When
  472. consumer, err := master.ConsumePartition("my_topic", 0, 1)
  473. if err != nil {
  474. t.Fatal(err)
  475. }
  476. assertMessageOffset(t, <-consumer.Messages(), 1)
  477. assertMessageOffset(t, <-consumer.Messages(), 2)
  478. safeClose(t, consumer)
  479. safeClose(t, master)
  480. broker0.Close()
  481. }
  482. func TestConsumeMessageWithSessionIDs(t *testing.T) {
  483. // Given
  484. fetchResponse1 := &FetchResponse{Version: 7}
  485. fetchResponse1.AddMessage("my_topic", 0, nil, testMsg, 1)
  486. fetchResponse1.AddMessage("my_topic", 0, nil, testMsg, 2)
  487. cfg := NewConfig()
  488. cfg.Version = V1_1_0_0
  489. broker0 := NewMockBroker(t, 0)
  490. fetchResponse2 := &FetchResponse{}
  491. fetchResponse2.Version = 7
  492. fetchResponse2.AddError("my_topic", 0, ErrNoError)
  493. broker0.SetHandlerByMap(map[string]MockResponse{
  494. "MetadataRequest": NewMockMetadataResponse(t).
  495. SetBroker(broker0.Addr(), broker0.BrokerID()).
  496. SetLeader("my_topic", 0, broker0.BrokerID()),
  497. "OffsetRequest": NewMockOffsetResponse(t).
  498. SetVersion(1).
  499. SetOffset("my_topic", 0, OffsetNewest, 1234).
  500. SetOffset("my_topic", 0, OffsetOldest, 0),
  501. "FetchRequest": NewMockSequence(fetchResponse1, fetchResponse2),
  502. })
  503. master, err := NewConsumer([]string{broker0.Addr()}, cfg)
  504. if err != nil {
  505. t.Fatal(err)
  506. }
  507. // When
  508. consumer, err := master.ConsumePartition("my_topic", 0, 1)
  509. if err != nil {
  510. t.Fatal(err)
  511. }
  512. assertMessageOffset(t, <-consumer.Messages(), 1)
  513. assertMessageOffset(t, <-consumer.Messages(), 2)
  514. safeClose(t, consumer)
  515. safeClose(t, master)
  516. broker0.Close()
  517. fetchReq := broker0.History()[3].Request.(*FetchRequest)
  518. if fetchReq.SessionID != 0 || fetchReq.Epoch != -1 {
  519. t.Error("Expected session ID to be zero & Epoch to be -1")
  520. }
  521. }
  522. // It is fine if offsets of fetched messages are not sequential (although
  523. // strictly increasing!).
  524. func TestConsumerNonSequentialOffsets(t *testing.T) {
  525. // Given
  526. legacyFetchResponse := &FetchResponse{}
  527. legacyFetchResponse.AddMessage("my_topic", 0, nil, testMsg, 5)
  528. legacyFetchResponse.AddMessage("my_topic", 0, nil, testMsg, 7)
  529. legacyFetchResponse.AddMessage("my_topic", 0, nil, testMsg, 11)
  530. newFetchResponse := &FetchResponse{Version: 4}
  531. newFetchResponse.AddRecord("my_topic", 0, nil, testMsg, 5)
  532. newFetchResponse.AddRecord("my_topic", 0, nil, testMsg, 7)
  533. newFetchResponse.AddRecord("my_topic", 0, nil, testMsg, 11)
  534. newFetchResponse.SetLastOffsetDelta("my_topic", 0, 11)
  535. newFetchResponse.SetLastStableOffset("my_topic", 0, 11)
  536. for _, fetchResponse1 := range []*FetchResponse{legacyFetchResponse, newFetchResponse} {
  537. var offsetResponseVersion int16
  538. cfg := NewConfig()
  539. if fetchResponse1.Version >= 4 {
  540. cfg.Version = V0_11_0_0
  541. offsetResponseVersion = 1
  542. }
  543. broker0 := NewMockBroker(t, 0)
  544. fetchResponse2 := &FetchResponse{Version: fetchResponse1.Version}
  545. fetchResponse2.AddError("my_topic", 0, ErrNoError)
  546. broker0.SetHandlerByMap(map[string]MockResponse{
  547. "MetadataRequest": NewMockMetadataResponse(t).
  548. SetBroker(broker0.Addr(), broker0.BrokerID()).
  549. SetLeader("my_topic", 0, broker0.BrokerID()),
  550. "OffsetRequest": NewMockOffsetResponse(t).
  551. SetVersion(offsetResponseVersion).
  552. SetOffset("my_topic", 0, OffsetNewest, 1234).
  553. SetOffset("my_topic", 0, OffsetOldest, 0),
  554. "FetchRequest": NewMockSequence(fetchResponse1, fetchResponse2),
  555. })
  556. master, err := NewConsumer([]string{broker0.Addr()}, cfg)
  557. if err != nil {
  558. t.Fatal(err)
  559. }
  560. // When
  561. consumer, err := master.ConsumePartition("my_topic", 0, 3)
  562. if err != nil {
  563. t.Fatal(err)
  564. }
  565. // Then: messages with offsets 1 and 2 are not returned even though they
  566. // are present in the response.
  567. assertMessageOffset(t, <-consumer.Messages(), 5)
  568. assertMessageOffset(t, <-consumer.Messages(), 7)
  569. assertMessageOffset(t, <-consumer.Messages(), 11)
  570. safeClose(t, consumer)
  571. safeClose(t, master)
  572. broker0.Close()
  573. }
  574. }
  575. // If leadership for a partition is changing then consumer resolves the new
  576. // leader and switches to it.
  577. func TestConsumerRebalancingMultiplePartitions(t *testing.T) {
  578. // initial setup
  579. seedBroker := NewMockBroker(t, 10)
  580. leader0 := NewMockBroker(t, 0)
  581. leader1 := NewMockBroker(t, 1)
  582. seedBroker.SetHandlerByMap(map[string]MockResponse{
  583. "MetadataRequest": NewMockMetadataResponse(t).
  584. SetBroker(leader0.Addr(), leader0.BrokerID()).
  585. SetBroker(leader1.Addr(), leader1.BrokerID()).
  586. SetBroker(seedBroker.Addr(), seedBroker.BrokerID()).
  587. SetLeader("my_topic", 0, leader0.BrokerID()).
  588. SetLeader("my_topic", 1, leader1.BrokerID()),
  589. })
  590. mockOffsetResponse1 := NewMockOffsetResponse(t).
  591. SetOffset("my_topic", 0, OffsetOldest, 0).
  592. SetOffset("my_topic", 0, OffsetNewest, 1000).
  593. SetOffset("my_topic", 1, OffsetOldest, 0).
  594. SetOffset("my_topic", 1, OffsetNewest, 1000)
  595. leader0.SetHandlerByMap(map[string]MockResponse{
  596. "OffsetRequest": mockOffsetResponse1,
  597. "FetchRequest": NewMockFetchResponse(t, 1),
  598. })
  599. leader1.SetHandlerByMap(map[string]MockResponse{
  600. "OffsetRequest": mockOffsetResponse1,
  601. "FetchRequest": NewMockFetchResponse(t, 1),
  602. })
  603. // launch test goroutines
  604. config := NewConfig()
  605. config.Consumer.Retry.Backoff = 50
  606. master, err := NewConsumer([]string{seedBroker.Addr()}, config)
  607. if err != nil {
  608. t.Fatal(err)
  609. }
  610. // we expect to end up (eventually) consuming exactly ten messages on each partition
  611. var wg sync.WaitGroup
  612. for i := int32(0); i < 2; i++ {
  613. consumer, err := master.ConsumePartition("my_topic", i, 0)
  614. if err != nil {
  615. t.Error(err)
  616. }
  617. go func(c PartitionConsumer) {
  618. for err := range c.Errors() {
  619. t.Error(err)
  620. }
  621. }(consumer)
  622. wg.Add(1)
  623. go func(partition int32, c PartitionConsumer) {
  624. for i := 0; i < 10; i++ {
  625. message := <-consumer.Messages()
  626. if message.Offset != int64(i) {
  627. t.Error("Incorrect message offset!", i, partition, message.Offset)
  628. }
  629. if message.Partition != partition {
  630. t.Error("Incorrect message partition!")
  631. }
  632. }
  633. safeClose(t, consumer)
  634. wg.Done()
  635. }(i, consumer)
  636. }
  637. time.Sleep(50 * time.Millisecond)
  638. Logger.Printf(" STAGE 1")
  639. // Stage 1:
  640. // * my_topic/0 -> leader0 serves 4 messages
  641. // * my_topic/1 -> leader1 serves 0 messages
  642. mockFetchResponse := NewMockFetchResponse(t, 1)
  643. for i := 0; i < 4; i++ {
  644. mockFetchResponse.SetMessage("my_topic", 0, int64(i), testMsg)
  645. }
  646. leader0.SetHandlerByMap(map[string]MockResponse{
  647. "FetchRequest": mockFetchResponse,
  648. })
  649. time.Sleep(50 * time.Millisecond)
  650. Logger.Printf(" STAGE 2")
  651. // Stage 2:
  652. // * leader0 says that it is no longer serving my_topic/0
  653. // * seedBroker tells that leader1 is serving my_topic/0 now
  654. // seed broker tells that the new partition 0 leader is leader1
  655. seedBroker.SetHandlerByMap(map[string]MockResponse{
  656. "MetadataRequest": NewMockMetadataResponse(t).
  657. SetLeader("my_topic", 0, leader1.BrokerID()).
  658. SetLeader("my_topic", 1, leader1.BrokerID()).
  659. SetBroker(leader0.Addr(), leader0.BrokerID()).
  660. SetBroker(leader1.Addr(), leader1.BrokerID()).
  661. SetBroker(seedBroker.Addr(), seedBroker.BrokerID()),
  662. })
  663. // leader0 says no longer leader of partition 0
  664. fetchResponse := new(FetchResponse)
  665. fetchResponse.AddError("my_topic", 0, ErrNotLeaderForPartition)
  666. leader0.SetHandlerByMap(map[string]MockResponse{
  667. "FetchRequest": NewMockWrapper(fetchResponse),
  668. })
  669. time.Sleep(50 * time.Millisecond)
  670. Logger.Printf(" STAGE 3")
  671. // Stage 3:
  672. // * my_topic/0 -> leader1 serves 3 messages
  673. // * my_topic/1 -> leader1 server 8 messages
  674. // leader1 provides 3 message on partition 0, and 8 messages on partition 1
  675. mockFetchResponse2 := NewMockFetchResponse(t, 2)
  676. for i := 4; i < 7; i++ {
  677. mockFetchResponse2.SetMessage("my_topic", 0, int64(i), testMsg)
  678. }
  679. for i := 0; i < 8; i++ {
  680. mockFetchResponse2.SetMessage("my_topic", 1, int64(i), testMsg)
  681. }
  682. leader1.SetHandlerByMap(map[string]MockResponse{
  683. "FetchRequest": mockFetchResponse2,
  684. })
  685. time.Sleep(50 * time.Millisecond)
  686. Logger.Printf(" STAGE 4")
  687. // Stage 4:
  688. // * my_topic/0 -> leader1 serves 3 messages
  689. // * my_topic/1 -> leader1 tells that it is no longer the leader
  690. // * seedBroker tells that leader0 is a new leader for my_topic/1
  691. // metadata assigns 0 to leader1 and 1 to leader0
  692. seedBroker.SetHandlerByMap(map[string]MockResponse{
  693. "MetadataRequest": NewMockMetadataResponse(t).
  694. SetLeader("my_topic", 0, leader1.BrokerID()).
  695. SetLeader("my_topic", 1, leader0.BrokerID()).
  696. SetBroker(leader0.Addr(), leader0.BrokerID()).
  697. SetBroker(leader1.Addr(), leader1.BrokerID()).
  698. SetBroker(seedBroker.Addr(), seedBroker.BrokerID()),
  699. })
  700. // leader1 provides three more messages on partition0, says no longer leader of partition1
  701. mockFetchResponse3 := NewMockFetchResponse(t, 3).
  702. SetMessage("my_topic", 0, int64(7), testMsg).
  703. SetMessage("my_topic", 0, int64(8), testMsg).
  704. SetMessage("my_topic", 0, int64(9), testMsg)
  705. fetchResponse4 := new(FetchResponse)
  706. fetchResponse4.AddError("my_topic", 1, ErrNotLeaderForPartition)
  707. leader1.SetHandlerByMap(map[string]MockResponse{
  708. "FetchRequest": NewMockSequence(mockFetchResponse3, fetchResponse4),
  709. })
  710. // leader0 provides two messages on partition 1
  711. mockFetchResponse4 := NewMockFetchResponse(t, 2)
  712. for i := 8; i < 10; i++ {
  713. mockFetchResponse4.SetMessage("my_topic", 1, int64(i), testMsg)
  714. }
  715. leader0.SetHandlerByMap(map[string]MockResponse{
  716. "FetchRequest": mockFetchResponse4,
  717. })
  718. wg.Wait()
  719. safeClose(t, master)
  720. leader1.Close()
  721. leader0.Close()
  722. seedBroker.Close()
  723. }
  724. // When two partitions have the same broker as the leader, if one partition
  725. // consumer channel buffer is full then that does not affect the ability to
  726. // read messages by the other consumer.
  727. func TestConsumerInterleavedClose(t *testing.T) {
  728. // Given
  729. broker0 := NewMockBroker(t, 0)
  730. broker0.SetHandlerByMap(map[string]MockResponse{
  731. "MetadataRequest": NewMockMetadataResponse(t).
  732. SetBroker(broker0.Addr(), broker0.BrokerID()).
  733. SetLeader("my_topic", 0, broker0.BrokerID()).
  734. SetLeader("my_topic", 1, broker0.BrokerID()),
  735. "OffsetRequest": NewMockOffsetResponse(t).
  736. SetOffset("my_topic", 0, OffsetOldest, 1000).
  737. SetOffset("my_topic", 0, OffsetNewest, 1100).
  738. SetOffset("my_topic", 1, OffsetOldest, 2000).
  739. SetOffset("my_topic", 1, OffsetNewest, 2100),
  740. "FetchRequest": NewMockFetchResponse(t, 1).
  741. SetMessage("my_topic", 0, 1000, testMsg).
  742. SetMessage("my_topic", 0, 1001, testMsg).
  743. SetMessage("my_topic", 0, 1002, testMsg).
  744. SetMessage("my_topic", 1, 2000, testMsg),
  745. })
  746. config := NewConfig()
  747. config.ChannelBufferSize = 0
  748. master, err := NewConsumer([]string{broker0.Addr()}, config)
  749. if err != nil {
  750. t.Fatal(err)
  751. }
  752. c0, err := master.ConsumePartition("my_topic", 0, 1000)
  753. if err != nil {
  754. t.Fatal(err)
  755. }
  756. c1, err := master.ConsumePartition("my_topic", 1, 2000)
  757. if err != nil {
  758. t.Fatal(err)
  759. }
  760. // When/Then: we can read from partition 0 even if nobody reads from partition 1
  761. assertMessageOffset(t, <-c0.Messages(), 1000)
  762. assertMessageOffset(t, <-c0.Messages(), 1001)
  763. assertMessageOffset(t, <-c0.Messages(), 1002)
  764. safeClose(t, c1)
  765. safeClose(t, c0)
  766. safeClose(t, master)
  767. broker0.Close()
  768. }
  769. func TestConsumerBounceWithReferenceOpen(t *testing.T) {
  770. broker0 := NewMockBroker(t, 0)
  771. broker0Addr := broker0.Addr()
  772. broker1 := NewMockBroker(t, 1)
  773. mockMetadataResponse := NewMockMetadataResponse(t).
  774. SetBroker(broker0.Addr(), broker0.BrokerID()).
  775. SetBroker(broker1.Addr(), broker1.BrokerID()).
  776. SetLeader("my_topic", 0, broker0.BrokerID()).
  777. SetLeader("my_topic", 1, broker1.BrokerID())
  778. mockOffsetResponse := NewMockOffsetResponse(t).
  779. SetOffset("my_topic", 0, OffsetOldest, 1000).
  780. SetOffset("my_topic", 0, OffsetNewest, 1100).
  781. SetOffset("my_topic", 1, OffsetOldest, 2000).
  782. SetOffset("my_topic", 1, OffsetNewest, 2100)
  783. mockFetchResponse := NewMockFetchResponse(t, 1)
  784. for i := 0; i < 10; i++ {
  785. mockFetchResponse.SetMessage("my_topic", 0, int64(1000+i), testMsg)
  786. mockFetchResponse.SetMessage("my_topic", 1, int64(2000+i), testMsg)
  787. }
  788. broker0.SetHandlerByMap(map[string]MockResponse{
  789. "OffsetRequest": mockOffsetResponse,
  790. "FetchRequest": mockFetchResponse,
  791. })
  792. broker1.SetHandlerByMap(map[string]MockResponse{
  793. "MetadataRequest": mockMetadataResponse,
  794. "OffsetRequest": mockOffsetResponse,
  795. "FetchRequest": mockFetchResponse,
  796. })
  797. config := NewConfig()
  798. config.Consumer.Return.Errors = true
  799. config.Consumer.Retry.Backoff = 100 * time.Millisecond
  800. config.ChannelBufferSize = 1
  801. master, err := NewConsumer([]string{broker1.Addr()}, config)
  802. if err != nil {
  803. t.Fatal(err)
  804. }
  805. c0, err := master.ConsumePartition("my_topic", 0, 1000)
  806. if err != nil {
  807. t.Fatal(err)
  808. }
  809. c1, err := master.ConsumePartition("my_topic", 1, 2000)
  810. if err != nil {
  811. t.Fatal(err)
  812. }
  813. // read messages from both partition to make sure that both brokers operate
  814. // normally.
  815. assertMessageOffset(t, <-c0.Messages(), 1000)
  816. assertMessageOffset(t, <-c1.Messages(), 2000)
  817. // Simulate broker shutdown. Note that metadata response does not change,
  818. // that is the leadership does not move to another broker. So partition
  819. // consumer will keep retrying to restore the connection with the broker.
  820. broker0.Close()
  821. // Make sure that while the partition/0 leader is down, consumer/partition/1
  822. // is capable of pulling messages from broker1.
  823. for i := 1; i < 7; i++ {
  824. offset := (<-c1.Messages()).Offset
  825. if offset != int64(2000+i) {
  826. t.Errorf("Expected offset %d from consumer/partition/1", int64(2000+i))
  827. }
  828. }
  829. // Bring broker0 back to service.
  830. broker0 = NewMockBrokerAddr(t, 0, broker0Addr)
  831. broker0.SetHandlerByMap(map[string]MockResponse{
  832. "FetchRequest": mockFetchResponse,
  833. })
  834. // Read the rest of messages from both partitions.
  835. for i := 7; i < 10; i++ {
  836. assertMessageOffset(t, <-c1.Messages(), int64(2000+i))
  837. }
  838. for i := 1; i < 10; i++ {
  839. assertMessageOffset(t, <-c0.Messages(), int64(1000+i))
  840. }
  841. select {
  842. case <-c0.Errors():
  843. default:
  844. t.Errorf("Partition consumer should have detected broker restart")
  845. }
  846. safeClose(t, c1)
  847. safeClose(t, c0)
  848. safeClose(t, master)
  849. broker0.Close()
  850. broker1.Close()
  851. }
  852. func TestConsumerOffsetOutOfRange(t *testing.T) {
  853. // Given
  854. broker0 := NewMockBroker(t, 2)
  855. broker0.SetHandlerByMap(map[string]MockResponse{
  856. "MetadataRequest": NewMockMetadataResponse(t).
  857. SetBroker(broker0.Addr(), broker0.BrokerID()).
  858. SetLeader("my_topic", 0, broker0.BrokerID()),
  859. "OffsetRequest": NewMockOffsetResponse(t).
  860. SetOffset("my_topic", 0, OffsetNewest, 1234).
  861. SetOffset("my_topic", 0, OffsetOldest, 2345),
  862. })
  863. master, err := NewConsumer([]string{broker0.Addr()}, nil)
  864. if err != nil {
  865. t.Fatal(err)
  866. }
  867. // When/Then
  868. if _, err := master.ConsumePartition("my_topic", 0, 0); err != ErrOffsetOutOfRange {
  869. t.Fatal("Should return ErrOffsetOutOfRange, got:", err)
  870. }
  871. if _, err := master.ConsumePartition("my_topic", 0, 3456); err != ErrOffsetOutOfRange {
  872. t.Fatal("Should return ErrOffsetOutOfRange, got:", err)
  873. }
  874. if _, err := master.ConsumePartition("my_topic", 0, -3); err != ErrOffsetOutOfRange {
  875. t.Fatal("Should return ErrOffsetOutOfRange, got:", err)
  876. }
  877. safeClose(t, master)
  878. broker0.Close()
  879. }
  880. func TestConsumerExpiryTicker(t *testing.T) {
  881. // Given
  882. broker0 := NewMockBroker(t, 0)
  883. fetchResponse1 := &FetchResponse{}
  884. for i := 1; i <= 8; i++ {
  885. fetchResponse1.AddMessage("my_topic", 0, nil, testMsg, int64(i))
  886. }
  887. broker0.SetHandlerByMap(map[string]MockResponse{
  888. "MetadataRequest": NewMockMetadataResponse(t).
  889. SetBroker(broker0.Addr(), broker0.BrokerID()).
  890. SetLeader("my_topic", 0, broker0.BrokerID()),
  891. "OffsetRequest": NewMockOffsetResponse(t).
  892. SetOffset("my_topic", 0, OffsetNewest, 1234).
  893. SetOffset("my_topic", 0, OffsetOldest, 1),
  894. "FetchRequest": NewMockSequence(fetchResponse1),
  895. })
  896. config := NewConfig()
  897. config.ChannelBufferSize = 0
  898. config.Consumer.MaxProcessingTime = 10 * time.Millisecond
  899. master, err := NewConsumer([]string{broker0.Addr()}, config)
  900. if err != nil {
  901. t.Fatal(err)
  902. }
  903. // When
  904. consumer, err := master.ConsumePartition("my_topic", 0, 1)
  905. if err != nil {
  906. t.Fatal(err)
  907. }
  908. // Then: messages with offsets 1 through 8 are read
  909. for i := 1; i <= 8; i++ {
  910. assertMessageOffset(t, <-consumer.Messages(), int64(i))
  911. time.Sleep(2 * time.Millisecond)
  912. }
  913. safeClose(t, consumer)
  914. safeClose(t, master)
  915. broker0.Close()
  916. }
  917. func TestConsumerTimestamps(t *testing.T) {
  918. now := time.Now().Truncate(time.Millisecond)
  919. type testMessage struct {
  920. key Encoder
  921. offset int64
  922. timestamp time.Time
  923. }
  924. for _, d := range []struct {
  925. kversion KafkaVersion
  926. logAppendTime bool
  927. messages []testMessage
  928. expectedTimestamp []time.Time
  929. }{
  930. {MinVersion, false, []testMessage{
  931. {testMsg, 1, now},
  932. {testMsg, 2, now},
  933. }, []time.Time{{}, {}}},
  934. {V0_9_0_0, false, []testMessage{
  935. {testMsg, 1, now},
  936. {testMsg, 2, now},
  937. }, []time.Time{{}, {}}},
  938. {V0_10_0_0, false, []testMessage{
  939. {testMsg, 1, now},
  940. {testMsg, 2, now},
  941. }, []time.Time{{}, {}}},
  942. {V0_10_2_1, false, []testMessage{
  943. {testMsg, 1, now.Add(time.Second)},
  944. {testMsg, 2, now.Add(2 * time.Second)},
  945. }, []time.Time{now.Add(time.Second), now.Add(2 * time.Second)}},
  946. {V0_10_2_1, true, []testMessage{
  947. {testMsg, 1, now.Add(time.Second)},
  948. {testMsg, 2, now.Add(2 * time.Second)},
  949. }, []time.Time{now, now}},
  950. {V0_11_0_0, false, []testMessage{
  951. {testMsg, 1, now.Add(time.Second)},
  952. {testMsg, 2, now.Add(2 * time.Second)},
  953. }, []time.Time{now.Add(time.Second), now.Add(2 * time.Second)}},
  954. {V0_11_0_0, true, []testMessage{
  955. {testMsg, 1, now.Add(time.Second)},
  956. {testMsg, 2, now.Add(2 * time.Second)},
  957. }, []time.Time{now, now}},
  958. } {
  959. var fr *FetchResponse
  960. var offsetResponseVersion int16
  961. cfg := NewConfig()
  962. cfg.Version = d.kversion
  963. switch {
  964. case d.kversion.IsAtLeast(V0_11_0_0):
  965. offsetResponseVersion = 1
  966. fr = &FetchResponse{Version: 4, LogAppendTime: d.logAppendTime, Timestamp: now}
  967. for _, m := range d.messages {
  968. fr.AddRecordWithTimestamp("my_topic", 0, m.key, testMsg, m.offset, m.timestamp)
  969. }
  970. fr.SetLastOffsetDelta("my_topic", 0, 2)
  971. fr.SetLastStableOffset("my_topic", 0, 2)
  972. case d.kversion.IsAtLeast(V0_10_1_0):
  973. offsetResponseVersion = 1
  974. fr = &FetchResponse{Version: 3, LogAppendTime: d.logAppendTime, Timestamp: now}
  975. for _, m := range d.messages {
  976. fr.AddMessageWithTimestamp("my_topic", 0, m.key, testMsg, m.offset, m.timestamp, 1)
  977. }
  978. default:
  979. var version int16
  980. switch {
  981. case d.kversion.IsAtLeast(V0_10_0_0):
  982. version = 2
  983. case d.kversion.IsAtLeast(V0_9_0_0):
  984. version = 1
  985. }
  986. fr = &FetchResponse{Version: version}
  987. for _, m := range d.messages {
  988. fr.AddMessageWithTimestamp("my_topic", 0, m.key, testMsg, m.offset, m.timestamp, 0)
  989. }
  990. }
  991. broker0 := NewMockBroker(t, 0)
  992. broker0.SetHandlerByMap(map[string]MockResponse{
  993. "MetadataRequest": NewMockMetadataResponse(t).
  994. SetBroker(broker0.Addr(), broker0.BrokerID()).
  995. SetLeader("my_topic", 0, broker0.BrokerID()),
  996. "OffsetRequest": NewMockOffsetResponse(t).
  997. SetVersion(offsetResponseVersion).
  998. SetOffset("my_topic", 0, OffsetNewest, 1234).
  999. SetOffset("my_topic", 0, OffsetOldest, 0),
  1000. "FetchRequest": NewMockSequence(fr),
  1001. })
  1002. master, err := NewConsumer([]string{broker0.Addr()}, cfg)
  1003. if err != nil {
  1004. t.Fatal(err)
  1005. }
  1006. consumer, err := master.ConsumePartition("my_topic", 0, 1)
  1007. if err != nil {
  1008. t.Fatal(err)
  1009. }
  1010. for i, ts := range d.expectedTimestamp {
  1011. select {
  1012. case msg := <-consumer.Messages():
  1013. assertMessageOffset(t, msg, int64(i)+1)
  1014. if msg.Timestamp != ts {
  1015. t.Errorf("Wrong timestamp (kversion:%v, logAppendTime:%v): got: %v, want: %v",
  1016. d.kversion, d.logAppendTime, msg.Timestamp, ts)
  1017. }
  1018. case err := <-consumer.Errors():
  1019. t.Fatal(err)
  1020. }
  1021. }
  1022. safeClose(t, consumer)
  1023. safeClose(t, master)
  1024. broker0.Close()
  1025. }
  1026. }
  1027. // When set to ReadCommitted, no uncommitted message should be available in messages channel
  1028. func TestExcludeUncommitted(t *testing.T) {
  1029. // Given
  1030. broker0 := NewMockBroker(t, 0)
  1031. fetchResponse := &FetchResponse{
  1032. Version: 4,
  1033. Blocks: map[string]map[int32]*FetchResponseBlock{"my_topic": {0: {
  1034. AbortedTransactions: []*AbortedTransaction{{ProducerID: 7, FirstOffset: 1235}},
  1035. }}},
  1036. }
  1037. fetchResponse.AddRecordBatch("my_topic", 0, nil, testMsg, 1234, 7, true) // committed msg
  1038. fetchResponse.AddRecordBatch("my_topic", 0, nil, testMsg, 1235, 7, true) // uncommitted msg
  1039. fetchResponse.AddRecordBatch("my_topic", 0, nil, testMsg, 1236, 7, true) // uncommitted msg
  1040. fetchResponse.AddControlRecord("my_topic", 0, 1237, 7, ControlRecordAbort) // abort control record
  1041. fetchResponse.AddRecordBatch("my_topic", 0, nil, testMsg, 1238, 7, true) // committed msg
  1042. broker0.SetHandlerByMap(map[string]MockResponse{
  1043. "MetadataRequest": NewMockMetadataResponse(t).
  1044. SetBroker(broker0.Addr(), broker0.BrokerID()).
  1045. SetLeader("my_topic", 0, broker0.BrokerID()),
  1046. "OffsetRequest": NewMockOffsetResponse(t).
  1047. SetVersion(1).
  1048. SetOffset("my_topic", 0, OffsetOldest, 0).
  1049. SetOffset("my_topic", 0, OffsetNewest, 1237),
  1050. "FetchRequest": NewMockWrapper(fetchResponse),
  1051. })
  1052. cfg := NewConfig()
  1053. cfg.Consumer.Return.Errors = true
  1054. cfg.Version = V0_11_0_0
  1055. cfg.Consumer.IsolationLevel = ReadCommitted
  1056. // When
  1057. master, err := NewConsumer([]string{broker0.Addr()}, cfg)
  1058. if err != nil {
  1059. t.Fatal(err)
  1060. }
  1061. consumer, err := master.ConsumePartition("my_topic", 0, 1234)
  1062. if err != nil {
  1063. t.Fatal(err)
  1064. }
  1065. // Then: only the 2 committed messages are returned
  1066. select {
  1067. case message := <-consumer.Messages():
  1068. assertMessageOffset(t, message, int64(1234))
  1069. case err := <-consumer.Errors():
  1070. t.Error(err)
  1071. }
  1072. select {
  1073. case message := <-consumer.Messages():
  1074. assertMessageOffset(t, message, int64(1238))
  1075. case err := <-consumer.Errors():
  1076. t.Error(err)
  1077. }
  1078. safeClose(t, consumer)
  1079. safeClose(t, master)
  1080. broker0.Close()
  1081. }
  1082. func assertMessageOffset(t *testing.T, msg *ConsumerMessage, expectedOffset int64) {
  1083. if msg.Offset != expectedOffset {
  1084. t.Errorf("Incorrect message offset: expected=%d, actual=%d", expectedOffset, msg.Offset)
  1085. }
  1086. }
  1087. // This example shows how to use the consumer to read messages
  1088. // from a single partition.
  1089. func ExampleConsumer() {
  1090. consumer, err := NewConsumer([]string{"localhost:9092"}, nil)
  1091. if err != nil {
  1092. panic(err)
  1093. }
  1094. defer func() {
  1095. if err := consumer.Close(); err != nil {
  1096. log.Fatalln(err)
  1097. }
  1098. }()
  1099. partitionConsumer, err := consumer.ConsumePartition("my_topic", 0, OffsetNewest)
  1100. if err != nil {
  1101. panic(err)
  1102. }
  1103. defer func() {
  1104. if err := partitionConsumer.Close(); err != nil {
  1105. log.Fatalln(err)
  1106. }
  1107. }()
  1108. // Trap SIGINT to trigger a shutdown.
  1109. signals := make(chan os.Signal, 1)
  1110. signal.Notify(signals, os.Interrupt)
  1111. consumed := 0
  1112. ConsumerLoop:
  1113. for {
  1114. select {
  1115. case msg := <-partitionConsumer.Messages():
  1116. log.Printf("Consumed message offset %d\n", msg.Offset)
  1117. consumed++
  1118. case <-signals:
  1119. break ConsumerLoop
  1120. }
  1121. }
  1122. log.Printf("Consumed: %d\n", consumed)
  1123. }
  1124. func Test_partitionConsumer_parseResponse(t *testing.T) {
  1125. type args struct {
  1126. response *FetchResponse
  1127. }
  1128. tests := []struct {
  1129. name string
  1130. args args
  1131. want []*ConsumerMessage
  1132. wantErr bool
  1133. }{
  1134. {
  1135. name: "empty but throttled FetchResponse is not considered an error",
  1136. args: args{
  1137. response: &FetchResponse{
  1138. ThrottleTime: time.Millisecond,
  1139. },
  1140. },
  1141. },
  1142. {
  1143. name: "empty FetchResponse is considered an incomplete response by default",
  1144. args: args{
  1145. response: &FetchResponse{},
  1146. },
  1147. wantErr: true,
  1148. },
  1149. }
  1150. for _, tt := range tests {
  1151. t.Run(tt.name, func(t *testing.T) {
  1152. child := &partitionConsumer{
  1153. broker: &brokerConsumer{
  1154. broker: &Broker{},
  1155. },
  1156. conf: &Config{},
  1157. }
  1158. got, err := child.parseResponse(tt.args.response)
  1159. if (err != nil) != tt.wantErr {
  1160. t.Errorf("partitionConsumer.parseResponse() error = %v, wantErr %v", err, tt.wantErr)
  1161. return
  1162. }
  1163. if !reflect.DeepEqual(got, tt.want) {
  1164. t.Errorf("partitionConsumer.parseResponse() = %v, want %v", got, tt.want)
  1165. }
  1166. })
  1167. }
  1168. }