Add credit-based flow control for stream consumers Stream consumers had no explicit flow control. During replay, broker delivery could outrun processing and overload downstream systems. Use creditWhenHalfMessagesProcessed(initialCredits) when building the stream consumer. Call context.processed() after each handled message so credits replenish as processing advances. Add stream.initialCredits to [stream] configuration. Its default is 1 and its minimum value is 1. Stream credits permit broker chunks, not individual messages, and chunk sizes vary. Queue consumers use prefetch 300, allowing up to 300 unacknowledged deliveries. Stream consumers use chunk-based credits, so the values are not directly comparable. The one-credit default keeps the initial stream window to one chunk while replenishment follows processing progress. Change-Id: I20f3050c8c6f52d98b82b0f86c5178ec77c53b95
diff --git a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/config/section/Stream.java b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/config/section/Stream.java index 7ac3059..8b4a980 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/config/section/Stream.java +++ b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/config/section/Stream.java
@@ -16,6 +16,7 @@ import com.google.common.flogger.FluentLogger; import com.googlesource.gerrit.plugins.rabbitmq.annotation.Default; +import com.googlesource.gerrit.plugins.rabbitmq.annotation.Limit; import java.util.ArrayList; import java.util.List; @@ -35,6 +36,10 @@ @Default("500") public Integer windowSize; + @Default("1") + @Limit(min = 1) + public Integer initialCredits; + public boolean isValid() { List<String> missingSettings = new ArrayList<>(3);
diff --git a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/session/type/StreamSubscriberSession.java b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/session/type/StreamSubscriberSession.java index 22806aa..c022349 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/session/type/StreamSubscriberSession.java +++ b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/session/type/StreamSubscriberSession.java
@@ -23,6 +23,7 @@ import com.googlesource.gerrit.plugins.rabbitmq.session.OffsetInfo; import com.googlesource.gerrit.plugins.rabbitmq.session.SubscriberSession; import com.rabbitmq.client.Channel; +import com.rabbitmq.stream.ConsumerFlowStrategy; import com.rabbitmq.stream.Message; import com.rabbitmq.stream.MessageHandler; import com.rabbitmq.stream.NoOffsetException; @@ -94,6 +95,10 @@ ctx.offsetSpecification(OffsetSpecification.offset(off + 1)); } }) + .flow() + .strategy( + ConsumerFlowStrategy.creditWhenHalfMessagesProcessed(streamProp.initialCredits)) + .builder() .manualTrackingStrategy() .builder() .listeners( @@ -289,6 +294,8 @@ } catch (IOException ex) { logger.atSevere().withCause(ex).log( "Error handling stream message with id %d", message.getPublishingId()); + } finally { + context.processed(); } }
diff --git a/src/main/resources/Documentation/config.md b/src/main/resources/Documentation/config.md index adbcc60..8934fc3 100644 --- a/src/main/resources/Documentation/config.md +++ b/src/main/resources/Documentation/config.md
@@ -94,7 +94,9 @@ in broker.config and is only used if `amqp.queuePrefix` is specified. * `amqp.consumerPrefetch` - * Decide how many events the client can queue for a consumer, defaults to 300. + * Decide how many events the client can queue for a queue-based consumer, defaults to 300. + This applies when `stream.enabled = false`. + When `stream.enabled = true`, stream consumer credits are used instead. * `exchange.name` * The name of exchange. @@ -123,6 +125,15 @@ offset of the currently proccessed message subtracted by `windowSize`, defaults to 500. Only used in broker.config. +* `stream.initialCredits` + * Initial credit window for stream consumers. Must be at least 1, defaults to 1. One credit + permits one broker chunk, not one event; the number of events in a chunk varies. Only used in + broker.config. This applies when `stream.enabled = true`. + RabbitMQ stream flow control documentation: + https://www.rabbitmq.com/docs/streams#flow-control + RabbitMQ stream Java client flow strategy API: + https://rabbitmq.github.io/rabbitmq-stream-java-client/stable/api/com/rabbitmq/stream/ConsumerFlowStrategy.html + * `general.publishAllGerritEvents` * Will publish gerrit stream events to configured exchange automatically if enabled, defaults to true.