From af6eb928a0dee8a758198dae343c60136bac1007 Mon Sep 17 00:00:00 2001 From: zmstone Date: Fri, 24 Oct 2025 12:44:53 +0200 Subject: [PATCH 1/3] chore: upgrade to kafka_protocol-4.2.9 --- changelog.md | 1 + rebar.config | 2 +- test/wolff_tests.erl | 4 ++-- 3 files changed, 4 insertions(+), 3 deletions(-) diff --git a/changelog.md b/changelog.md index c9c22c7..efe1825 100644 --- a/changelog.md +++ b/changelog.md @@ -1,6 +1,7 @@ * 4.1.0 - Fix 'failed' telemetry counter double-increment due to race condition. [#102](https://github.com/kafka4beam/wolff/pull/102) - Added producer process label `{wolff_producer, KafkaClientId, Topic, Partition}` [#102](https://github.com/kafka4beam/wolff/pull/102) + - Upgrade to `kafka_protocol-4.2.9` for better CRC32C performance. * 4.0.13 (merge 1.5.19) - Handle `record_list_too_large` error returned from Kafka. diff --git a/rebar.config b/rebar.config index 482ddb9..d6c1ce8 100644 --- a/rebar.config +++ b/rebar.config @@ -1,4 +1,4 @@ -{deps, [ {kafka_protocol, "4.2.8"} +{deps, [ {kafka_protocol, "4.2.9"} , {replayq, "0.4.1"} , {lc, "0.3.5"} , {telemetry, "1.1.0"} diff --git a/test/wolff_tests.erl b/test/wolff_tests.erl index 6b3608a..4408b2b 100644 --- a/test/wolff_tests.erl +++ b/test/wolff_tests.erl @@ -1030,8 +1030,8 @@ batch_bytes(Batch) -> wolff_producer:batch_bytes(Batch). encoded_bytes(Batch) -> - Encoded = kpro_batch:encode(2, Batch, no_compression), - iolist_size(Encoded). + {Bytes, _Encoded} = kpro_batch:encode(2, Batch, no_compression), + Bytes. create_topic(Topic, Partitions, ReplicationFactor, MaxMessageBytes, SegmentBytes) -> Cmd = create_topic_cmd(Topic, Partitions, ReplicationFactor, MaxMessageBytes, SegmentBytes), From f8b4a68821c4474a112e6e0f269c7f4e688566bd Mon Sep 17 00:00:00 2001 From: zmstone Date: Fri, 24 Oct 2025 20:36:26 +0200 Subject: [PATCH 2/3] refactor: make and send request in one func --- src/wolff_producer.erl | 27 +++++++++++++++------------ 1 file changed, 15 insertions(+), 12 deletions(-) diff --git a/src/wolff_producer.erl b/src/wolff_producer.erl index 82c4f05..892018c 100644 --- a/src/wolff_producer.erl +++ b/src/wolff_producer.erl @@ -423,27 +423,32 @@ use_defaults(Config, [{K, V} | Rest]) -> false -> use_defaults(Config#{K => V}, Rest) end. -remake_requests(#{q_items := Items} = Sent, St, ?reconnect) -> - #kpro_req{ref = Ref} = Req = make_request(Items, St), - {[Req], [Sent#{req_ref := Ref}]}; -remake_requests(#{q_items := Items} = Sent, St, ?message_too_large) -> - Reqs = lists:map(fun(Item) -> make_request([Item], St) end, Items), +resend_requests(#{q_items := Items} = Sent, St, ?reconnect) -> + {ok, Ref} = send_request(Items, St), + [Sent#{req_ref := Ref}]; +resend_requests(#{q_items := Items} = Sent, St, ?message_too_large) -> + Refs = lists:map(fun(Item) -> {ok, Ref} = send_request([Item], St), Ref end, Items), #{attempts := Attempts, q_ack_ref := Q_AckRef} = Sent, - {Reqs, make_sent_items_list(Items, Reqs, Attempts, Q_AckRef)}. + make_sent_items_list(Items, Refs, Attempts, Q_AckRef). %% only ack replayq when the last message is accepted by Kafka -make_sent_items_list([LastItem], [#kpro_req{ref = Ref}], Attempts, Q_AckRef) -> +make_sent_items_list([LastItem], [Ref], Attempts, Q_AckRef) -> [#{req_ref => Ref, q_items => [LastItem], q_ack_ref => Q_AckRef, attempts => Attempts }]; -make_sent_items_list([Item | Items], [#kpro_req{ref = Ref} | Reqs], Attempts, Q_AckRef) -> +make_sent_items_list([Item | Items], [Ref | Refs], Attempts, Q_AckRef) -> [#{req_ref => Ref, q_items => [Item], q_ack_ref => ?no_queue_ack, attempts => Attempts - } | make_sent_items_list(Items, Reqs, Attempts, Q_AckRef)]. + } | make_sent_items_list(Items, Refs, Attempts, Q_AckRef)]. + +send_request(QueueItems, #{conn := Conn} = St) -> + #kpro_req{ref = Ref} = Req = make_request(QueueItems, St), + ok = request_async(Conn, Req), + {ok, Ref}. make_request(QueueItems, #{config := #{required_acks := RequiredAcks, @@ -461,13 +466,11 @@ make_request(QueueItems, }). resend_sent_reqs(#{sent_reqs := SentReqs, - conn := Conn, config := Config } = St, Reason) -> F = fun(Sent, Acc) -> wolff_metrics:retried_inc(Config, 1), - {Reqs, NewSentList} = remake_requests(Sent, St, Reason), - lists:foreach(fun(Req) -> ok = request_async(Conn, Req) end, Reqs), + NewSentList = resend_requests(Sent, St, Reason), lists:reverse(NewSentList, Acc) end, NewSentReqs = lists:foldl(F, [], queue:to_list(SentReqs)), From cbde0da743b80bb662171d3cd38e51c10b304dfa Mon Sep 17 00:00:00 2001 From: zmstone Date: Fri, 24 Oct 2025 20:49:48 +0200 Subject: [PATCH 3/3] refactor: if message_too_large estimate batch instead of request size --- src/wolff_producer.erl | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/wolff_producer.erl b/src/wolff_producer.erl index 892018c..b15bb2a 100644 --- a/src/wolff_producer.erl +++ b/src/wolff_producer.erl @@ -623,8 +623,8 @@ do_handle_kafka_ack(?no_error, BaseOffset, St) -> do_handle_kafka_ack(EC, _BaseOffset, St) when EC =:= ?message_too_large orelse EC =:= ?record_list_too_large -> #{topic := Topic, partition := Partition, config := Config, sent_reqs := SentReqs} = St, {value, #{q_items := Items}} = queue:peek(SentReqs), - #kpro_req{msg = IoData} = make_request(Items, St), - Bytes = iolist_size(IoData), + Batch = get_batch_from_queue_items(Items), + Bytes = batch_bytes(Batch), Hint = case EC of ?message_too_large -> "Consider increasing server side topic config 'max.message.bytes'!"; @@ -638,7 +638,7 @@ do_handle_kafka_ack(EC, _BaseOffset, St) when EC =:= ?message_too_large orelse E 1 -> %% This is a single message batch, but it's still too large Note = "A single-request batch is dropped because it's too large for this topic! " ++ Hint, - log_error(Topic, Partition, EC, #{note => Note, encode_bytes => Bytes}), + log_error(Topic, Partition, EC, #{note => Note, estimated_batch_bytes => Bytes}), wolff_metrics:dropped_inc(Config, 1), clear_sent_and_ack_callers(EC, Reason, St); N ->