From 6044d83f97db34232f0be1fc5c39c97cc88cd515 Mon Sep 17 00:00:00 2001 From: Yiming Zang Date: Sat, 24 Jun 2023 21:23:56 -0700 Subject: [PATCH] [transactions] Throttle clients on certain error during handleInitProducerId --- .../streamnative/pulsar/handlers/kop/KafkaRequestHandler.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index f4de60ce0a..1252cd2174 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -2151,6 +2151,10 @@ protected void handleInitProducerId(KafkaHeaderAndRequest kafkaHeaderAndRequest, .setErrorCode(resp.getError().code()) .setProducerId(resp.getProducerId()) .setProducerEpoch(resp.getProducerEpoch()); + if (resp.getError() == Errors.COORDINATOR_LOAD_IN_PROGRESS + || resp.getError() == Errors.CONCURRENT_TRANSACTIONS) { + responseData.setThrottleTimeMs(1000); + } response.complete(new InitProducerIdResponse(responseData)); }); }