diff --git a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/SleuthStreamProperties.java b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/SleuthStreamProperties.java index 3ac65a286..0422de8e3 100644 --- a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/SleuthStreamProperties.java +++ b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/SleuthStreamProperties.java @@ -27,4 +27,5 @@ import lombok.Data; @Data public class SleuthStreamProperties { private boolean enabled = true; + private String group = SleuthSink.INPUT; } diff --git a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessor.java b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessor.java index 094ed9998..1ba4ead0a 100644 --- a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessor.java +++ b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamEnvironmentPostProcessor.java @@ -52,7 +52,8 @@ public class StreamEnvironmentPostProcessor implements EnvironmentPostProcessor SpringApplication application) { Map map = new HashMap(); ResourceLoader resourceLoader = application.getResourceLoader(); - resourceLoader = resourceLoader==null ? new DefaultResourceLoader() : resourceLoader; + resourceLoader = resourceLoader == null ? new DefaultResourceLoader() + : resourceLoader; PathMatchingResourcePatternResolver resolver = new PathMatchingResourcePatternResolver( resourceLoader); try { @@ -66,6 +67,11 @@ public class StreamEnvironmentPostProcessor implements EnvironmentPostProcessor catch (IOException e) { throw new IllegalStateException("Cannot load META-INF/spring.binders", e); } + // Technically this is only needed on the consumer, but it's fine to be explicit + // on producers as well. It puts all consumers in the same "group", meaning they + // compete with each other and only one gets each message. + map.put("spring.cloud.stream.bindings." + SleuthSink.INPUT + ".group", + environment.getProperty("spring.sleuth.stream.group", SleuthSink.INPUT)); addOrReplace(environment.getPropertySources(), map); }