From 4783b07407554bd3d5a727ef18fe498623b85acc Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Mon, 12 Jun 2017 15:43:11 +0100 Subject: [PATCH] Support for multi-valued (batched) events in stream aggregator This allows clients to send batched up events in the same format as before (or to continue to send single events). We can switch the default format to an array in 1.4.x. --- .../stream/HystrixStreamAggregator.java | 28 +++++++++++++++---- .../stream/HystrixStreamAggregatorTests.java | 26 +++++++++++------ 2 files changed, 40 insertions(+), 14 deletions(-) diff --git a/spring-cloud-netflix-turbine-stream/src/main/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregator.java b/spring-cloud-netflix-turbine-stream/src/main/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregator.java index ac50a54a..ca5b5f42 100644 --- a/spring-cloud-netflix-turbine-stream/src/main/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregator.java +++ b/spring-cloud-netflix-turbine-stream/src/main/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregator.java @@ -17,6 +17,7 @@ package org.springframework.cloud.netflix.turbine.stream; import java.io.IOException; +import java.util.List; import java.util.Map; import org.springframework.beans.factory.annotation.Autowired; @@ -56,18 +57,33 @@ public class HystrixStreamAggregator { payload = payload.replace("\\\"", "\""); } try { - @SuppressWarnings("unchecked") - Map map = this.objectMapper.readValue(payload, Map.class); - Map data = getPayloadData(map); - - log.debug("Received hystrix stream payload: " + data); - this.subject.onNext(data); + if (payload.startsWith("[")) { + @SuppressWarnings("unchecked") + List> list = this.objectMapper.readValue(payload, + List.class); + for (Map map : list) { + sendMap(map); + } + } + else { + @SuppressWarnings("unchecked") + Map map = this.objectMapper.readValue(payload, Map.class); + sendMap(map); + } } catch (IOException ex) { log.error("Error receiving hystrix stream payload: " + payload, ex); } } + private void sendMap(Map map) { + Map data = getPayloadData(map); + if (log.isDebugEnabled()) { + log.debug("Received hystrix stream payload: " + data); + } + this.subject.onNext(data); + } + public static Map getPayloadData(Map jsonMap) { @SuppressWarnings("unchecked") Map origin = (Map) jsonMap.get("origin"); diff --git a/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregatorTests.java b/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregatorTests.java index 998bb1ca..5d9352c3 100644 --- a/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregatorTests.java +++ b/spring-cloud-netflix-turbine-stream/src/test/java/org/springframework/cloud/netflix/turbine/stream/HystrixStreamAggregatorTests.java @@ -16,19 +16,20 @@ package org.springframework.cloud.netflix.turbine.stream; +import java.util.Map; + +import com.fasterxml.jackson.databind.ObjectMapper; + +import org.junit.Rule; +import org.junit.Test; + +import org.springframework.boot.test.rule.OutputCapture; + import static org.hamcrest.CoreMatchers.containsString; import static org.hamcrest.CoreMatchers.equalTo; import static org.hamcrest.CoreMatchers.not; import static org.junit.Assert.assertThat; -import java.util.Map; - -import org.junit.Rule; -import org.junit.Test; -import org.springframework.boot.test.rule.OutputCapture; - -import com.fasterxml.jackson.databind.ObjectMapper; - import rx.subjects.PublishSubject; public class HystrixStreamAggregatorTests { @@ -52,6 +53,15 @@ public class HystrixStreamAggregatorTests { this.output.expect(not(containsString("ERROR"))); } + @Test + public void messageWrappedInArray() throws Exception { + this.publisher.subscribe(map -> { + assertThat(map.get("type"), equalTo("HystrixCommand")); + }); + this.aggregator.sendToSubject("[" + PAYLOAD + "]"); + this.output.expect(not(containsString("ERROR"))); + } + @Test public void doubleEncodedMessage() throws Exception { this.publisher.subscribe(map -> {