Packages
brod
3.7.7
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
+19
additions
-12
deletions
| @@ -1,7 +1,7 @@ | |
| 1 1 | KAFKA_VERSION ?= 1.1 |
| 2 2 | PROJECT = brod |
| 3 3 | PROJECT_DESCRIPTION = Kafka client library in Erlang |
| 4 | - PROJECT_VERSION = 3.7.6 |
| 4 | + PROJECT_VERSION = 3.7.7 |
| 5 5 | |
| 6 6 | DEPS = supervisor3 kafka_protocol |
| @@ -1,5 +1,5 @@ | |
| 1 1 | {<<"name">>,<<"brod">>}. |
| 2 | - {<<"version">>,<<"3.7.6">>}. |
| 2 | + {<<"version">>,<<"3.7.7">>}. |
| 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.7.6"}, |
| 3 | + {vsn,"3.7.7"}, |
| 4 4 | {registered,[]}, |
| 5 5 | {applications,[kernel,stdlib,kafka_protocol,supervisor3]}, |
| 6 6 | {env,[]}, |
| @@ -473,19 +473,26 @@ handle_consumer_delivery(#kafka_message_set{ topic = Topic | |
| 473 473 | , partition = Partition |
| 474 474 | , messages = Messages |
| 475 475 | } = MsgSet, |
| 476 | - #state{message_type = MsgType} = State0) -> |
| 477 | - State = update_last_offset({Topic, Partition}, Messages, State0), |
| 478 | - case MsgType of |
| 479 | - message -> handle_messages(Topic, Partition, Messages, State); |
| 480 | - message_set -> handle_message_set(MsgSet, State) |
| 476 | + #state{ message_type = MsgType |
| 477 | + , consumers = Consumers0 |
| 478 | + } = State0) -> |
| 479 | + case get_consumer({Topic, Partition}, Consumers0) of |
| 480 | + #consumer{} = C -> |
| 481 | + Consumers = update_last_offset(Messages, C, Consumers0), |
| 482 | + State = State0#state{consumers = Consumers}, |
| 483 | + case MsgType of |
| 484 | + message -> handle_messages(Topic, Partition, Messages, State); |
| 485 | + message_set -> handle_message_set(MsgSet, State) |
| 486 | + end; |
| 487 | + false -> |
| 488 | + State0 |
| 481 489 | end. |
| 482 490 | |
| 483 | - update_last_offset(TP, Messages, #state{consumers = Consumers} = State) -> |
| 491 | + update_last_offset(Messages, Consumer0, Consumers) -> |
| 484 492 | %% brod_consumer never delivers empty message set, lists:last is safe |
| 485 493 | #kafka_message{offset = LastOffset} = lists:last(Messages), |
| 486 | - C = get_consumer(TP, Consumers), |
| 487 | - Consumer = C#consumer{last_offset = LastOffset}, |
| 488 | - State#state{consumers = put_consumer(Consumer, Consumers)}. |
| 494 | + Consumer = Consumer0#consumer{last_offset = LastOffset}, |
| 495 | + put_consumer(Consumer, Consumers). |
| 489 496 | |
| 490 497 | -spec start_subscribe_timer(?undef | reference(), timeout()) -> reference(). |
| 491 498 | start_subscribe_timer(?undef, Delay) -> |