Current section

113 Versions

Jump to

Compare versions

4 files changed
+21 additions
-5 deletions
  @@ -1,6 +1,6 @@
1 1 PROJECT = brod
2 2 PROJECT_DESCRIPTION = Kafka client library in Erlang
3 - PROJECT_VERSION = 3.5.1
3 + PROJECT_VERSION = 3.5.2
4 4
5 5 DEPS = supervisor3 kafka_protocol
6 6 TEST_DEPS = docopt jsone meck proper
  @@ -1,5 +1,5 @@
1 1 {<<"name">>,<<"brod">>}.
2 - {<<"version">>,<<"3.5.1">>}.
2 + {<<"version">>,<<"3.5.2">>}.
3 3 {<<"requirements">>,
4 4 #{<<"kafka_protocol">> =>
5 5 #{<<"app">> => <<"kafka_protocol">>,<<"optional">> => false,
  @@ -1,6 +1,6 @@
1 1 {application,brod,
2 2 [{description,"Apache Kafka Erlang client library"},
3 - {vsn,"3.5.1"},
3 + {vsn,"3.5.2"},
4 4 {registered,[]},
5 5 {applications,[kernel,stdlib,ssl,kafka_protocol,supervisor3]},
6 6 {env,[]},
  @@ -376,8 +376,24 @@ handle_message_set(#kafka_message_set{messages = ?incomplete_message(Size)},
376 376 State1 = State0#state{max_bytes = Size},
377 377 State = maybe_send_fetch_request(State1),
378 378 {noreply, State};
379 - handle_message_set(#kafka_message_set{messages = []}, State0) ->
380 - State = maybe_delay_fetch_request(State0),
379 + handle_message_set(#kafka_message_set{messages = [],
380 + high_wm_offset = HmOffset
381 + },
382 + #state{begin_offset = BeginOffset} = State0) ->
383 + State =
384 + case BeginOffset < HmOffset of
385 + true ->
386 + %% There are chances that kafka may return empty message set
387 + %% when messages are delete from a compacted topic.
388 + %% Since there is no way to know how big the 'hole' is
389 + %% we can only bump begin_offset with +1 and try again.
390 + State1 = State0#state{begin_offset = BeginOffset + 1},
391 + maybe_send_fetch_request(State1);
392 + false ->
393 + %% we have reached the end of partition
394 + %% try to poll again (maybe after a delay)
395 + maybe_delay_fetch_request(State0)
396 + end,
381 397 {noreply, State};
382 398 handle_message_set(#kafka_message_set{messages = Messages} = MsgSet,
383 399 #state{ subscriber = Subscriber