Current section

113 Versions

Jump to

Compare versions

5 files changed
+16 additions
-14 deletions
  @@ -1,14 +1,14 @@
1 1 KAFKA_VERSION ?= 1.1
2 2 PROJECT = brod
3 3 PROJECT_DESCRIPTION = Kafka client library in Erlang
4 - PROJECT_VERSION = 3.7.2
4 + PROJECT_VERSION = 3.7.3
5 5
6 6 DEPS = supervisor3 kafka_protocol
7 7
8 8 ERLC_OPTS = -Werror +warn_unused_vars +warn_shadow_vars +warn_unused_import +warn_obsolete_guard +debug_info -Dbuild_brod_cli
9 9
10 - dep_supervisor3_commit = 1.1.7
11 - dep_kafka_protocol_commit = 2.2.2
10 + dep_supervisor3_commit = 1.1.8
11 + dep_kafka_protocol_commit = 2.2.4
12 12 dep_kafka_protocol = git https://github.com/klarna/kafka_protocol.git $(dep_kafka_protocol_commit)
13 13
14 14 EDOC_OPTS = preprocess, {macros, [{build_brod_cli, true}]}
  @@ -1,12 +1,12 @@
1 1 {<<"name">>,<<"brod">>}.
2 - {<<"version">>,<<"3.7.2">>}.
2 + {<<"version">>,<<"3.7.3">>}.
3 3 {<<"requirements">>,
4 4 #{<<"kafka_protocol">> =>
5 5 #{<<"app">> => <<"kafka_protocol">>,<<"optional">> => false,
6 - <<"requirement">> => <<"2.2.2">>},
6 + <<"requirement">> => <<"2.2.4">>},
7 7 <<"supervisor3">> =>
8 8 #{<<"app">> => <<"supervisor3">>,<<"optional">> => false,
9 - <<"requirement">> => <<"1.1.7">>}}}.
9 + <<"requirement">> => <<"1.1.8">>}}}.
10 10 {<<"app">>,<<"brod">>}.
11 11 {<<"maintainers">>,[<<"Ivan Dyachkov">>,<<"Zaiming Shi">>]}.
12 12 {<<"precompiled">>,false}.
  @@ -1,5 +1,5 @@
1 - {deps, [ {supervisor3, "1.1.7"}
2 - , {kafka_protocol, "2.2.2"}
1 + {deps, [ {supervisor3, "1.1.8"}
2 + , {kafka_protocol, "2.2.4"}
3 3 ]}.
4 4 {edoc_opts, [{preprocess, true}, {macros, [{build_brod_cli, true}]}]}.
5 5 {erl_opts, [warn_unused_vars,warn_shadow_vars,warn_unused_import,warn_obsolete_guard,debug_info]}.
  @@ -1,6 +1,6 @@
1 1 {application,brod,
2 2 [{description,"Apache Kafka Erlang client library"},
3 - {vsn,"3.7.2"},
3 + {vsn,"3.7.3"},
4 4 {registered,[]},
5 5 {applications,[kernel,stdlib,kafka_protocol,supervisor3]},
6 6 {env,[]},
  @@ -552,16 +552,18 @@ handle_ack(AckRef, #state{ generationId = GenerationId
552 552 , coordinator = Coordinator
553 553 } = State, CommitNow) ->
554 554 {Topic, Partition, Offset} = AckRef,
555 - Consumer = get_consumer({Topic, Partition}, Consumers),
556 - #consumer{consumer_pid = ConsumerPid} = Consumer,
557 - ok = consume_ack(ConsumerPid, Offset),
558 - case CommitNow of
559 - true ->
555 + case get_consumer({Topic, Partition}, Consumers) of
556 + #consumer{consumer_pid = ConsumerPid} = Consumer when CommitNow ->
557 + ok = consume_ack(ConsumerPid, Offset),
560 558 ok = do_commit_ack(Coordinator, GenerationId, Topic, Partition, Offset),
561 559 NewConsumer = Consumer#consumer{acked_offset = Offset},
562 560 NewConsumers = put_consumer(NewConsumer, Consumers),
563 561 State#state{consumers = NewConsumers};
562 + #consumer{consumer_pid = ConsumerPid} ->
563 + ok = consume_ack(ConsumerPid, Offset),
564 + State;
564 565 false ->
566 + %% Stale async-ack, discard.
565 567 State
566 568 end.