Packages
brod
3.5.2
4.5.7
4.5.6
4.5.5
4.5.4
4.5.3
4.5.2
4.5.1
4.5.0
4.4.7
4.4.6
4.4.5
4.4.4
4.4.3
4.4.2
4.4.1
4.4.0
4.3.3
4.3.2
4.3.1
4.3.0
4.2.0
4.1.1
4.1.0
4.0.0
3.19.1
3.19.0
3.18.0
3.17.1
3.17.0
3.16.5
3.16.4
3.16.3
3.16.2
3.16.1
3.16.0
3.15.6
3.15.5
3.15.4
3.15.3
3.15.1
3.15.0
3.14.0
3.13.0
3.12.0
3.11.0
3.10.0
3.9.5
3.9.3
3.9.2
3.9.1
3.9.0
3.8.1
3.8.0
3.7.11
3.7.10
3.7.9
3.7.8
3.7.7
3.7.6
3.7.5
3.7.4
3.7.3
3.7.2
3.7.1
3.7.0
3.6.2
3.6.1
3.6.0
3.5.2
3.5.1
3.5.0
3.4.0
3.3.5
3.3.4
3.3.3
3.3.2
3.3.1
3.3.0
3.2.0
3.0.0
2.5.0
2.4.1
2.4.0
2.3.7
2.3.6
2.3.5
2.3.4
2.3.3
2.3.1
2.2.16
2.2.15
2.2.14
2.2.12
2.2.11
2.2.10
2.2.9
2.2.8
2.2.7
2.2.6
2.2.5
2.2.4
2.2.3
2.2.2
2.2.1
2.2.0
2.1.12
2.1.11
2.1.10
2.1.8
2.1.7
2.1.4
2.1.2
2.0.0
Apache Kafka Erlang client library
Current section
113 Versions
Jump to
Current section
113 Versions
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 |