From f966b4525cf5b5fead5e7650d5163da862cfb81a Mon Sep 17 00:00:00 2001 From: Spencer Gibb Date: Mon, 14 Jul 2014 20:13:45 -0600 Subject: [PATCH] use eureka to lookup hystrix streams in turbine --- README.md | 3 +- RUNNING.md | 2 +- pom.xml | 6 + spring-platform-netflix-turbine/pom.xml | 4 + .../platform/netflix/turbine/Application.java | 2 + .../turbine/EurekaInstanceDiscovery.java | 193 ++++++++++++++++++ .../turbine/SpringAggregatorFactory.java | 121 +++++++++++ .../netflix/turbine/SpringClusterMonitor.java | 79 +++++++ 8 files changed, 408 insertions(+), 2 deletions(-) create mode 100644 spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/EurekaInstanceDiscovery.java create mode 100644 spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/SpringAggregatorFactory.java create mode 100644 spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/SpringClusterMonitor.java diff --git a/README.md b/README.md index e33ded11..c4d76020 100644 --- a/README.md +++ b/README.md @@ -58,4 +58,5 @@ The default application name, virtual host and non-secure port are taken from th - [ ] Better observable example - [ ] Distributed refresh environment via platform bus - [x] Metrics aggregation (turbine) - - [ ] Use Eureka for instance discovery rather than static list see https://github.com/Netflix/Turbine/blob/master/turbine-contrib/src/main/java/com/netflix/turbine/discovery/EurekaInstanceDiscovery.java + - [x] Use Eureka for instance discovery rather than static list see https://github.com/Netflix/Turbine/blob/master/turbine-contrib/src/main/java/com/netflix/turbine/discovery/EurekaInstanceDiscovery.java + - [ ] Configure InstanceDiscovery.impl using auto config/config props diff --git a/RUNNING.md b/RUNNING.md index 6b1e1b85..da8dae4b 100644 --- a/RUNNING.md +++ b/RUNNING.md @@ -61,7 +61,7 @@ ## Netflix Turbine - `spring-platform-netflix-turbine$ java -jar target/spring-platform-netflix-turbine-1.0.0.BUILD-SNAPSHOT.war --turbine.aggregator.clusterConfig=sampleApps --turbine.ConfigPropertyBasedDiscovery.sampleApps.instances=localhost:9080 --turbine.instanceUrlSuffix=/hystrix.stream` + `spring-platform-netflix-turbine$ java -jar target/spring-platform-netflix-turbine-1.0.0.BUILD-SNAPSHOT.war --InstanceDiscovery.impl=io.spring.platform.netflix.turbine.EurekaInstanceDiscovery --turbine.appConfig=samplefrontendservice --turbine.aggregator.clusterConfig=sampleApps --turbine.instanceUrlSuffix=/hystrix.stream --turbine.instanceInsertPort=true` ## Sandbox Sample Backend diff --git a/pom.xml b/pom.xml index b94490ee..8320fc8c 100644 --- a/pom.xml +++ b/pom.xml @@ -47,6 +47,12 @@ com.netflix.eureka eureka-client ${eureka.version} + + + javax.servlet + servlet-api + + com.netflix.eureka diff --git a/spring-platform-netflix-turbine/pom.xml b/spring-platform-netflix-turbine/pom.xml index 90959c4c..100417b1 100644 --- a/spring-platform-netflix-turbine/pom.xml +++ b/spring-platform-netflix-turbine/pom.xml @@ -53,6 +53,10 @@ org.springframework.platform spring-platform-config-client + + com.netflix.eureka + eureka-client + com.netflix.turbine turbine-core diff --git a/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/Application.java b/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/Application.java index 9dc67f21..83292ec0 100644 --- a/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/Application.java +++ b/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/Application.java @@ -1,6 +1,7 @@ package io.spring.platform.netflix.turbine; import com.netflix.turbine.init.TurbineInit; +import com.netflix.turbine.plugins.PluginsFactory; import com.netflix.turbine.streaming.servlet.TurbineStreamServlet; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; @@ -47,6 +48,7 @@ public class Application extends SpringBootServletInitializer implements SmartLi @Override public void start() { + PluginsFactory.setClusterMonitorFactory(new SpringAggregatorFactory()); TurbineInit.init(); } diff --git a/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/EurekaInstanceDiscovery.java b/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/EurekaInstanceDiscovery.java new file mode 100644 index 00000000..f76ecac6 --- /dev/null +++ b/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/EurekaInstanceDiscovery.java @@ -0,0 +1,193 @@ +package io.spring.platform.netflix.turbine; + +import com.netflix.appinfo.AmazonInfo; +import com.netflix.appinfo.DataCenterInfo; +import com.netflix.appinfo.InstanceInfo; +import com.netflix.appinfo.InstanceInfo.InstanceStatus; +import com.netflix.config.DynamicPropertyFactory; +import com.netflix.config.DynamicStringProperty; +import com.netflix.discovery.DiscoveryManager; +import com.netflix.discovery.shared.Application; +import com.netflix.turbine.discovery.Instance; +import com.netflix.turbine.discovery.InstanceDiscovery; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.*; + +/** + * Class that encapsulates an {@link InstanceDiscovery} implementation that uses Eureka (see https://github.com/Netflix/eureka) + * The plugin requires a list of applications configured. It then queries the set of instances for each application. + * Instance information retrieved from Eureka must be translated to something that Turbine can understand i.e the {@link Instance} class. + * + * All the logic to perform this translation can be overriden here, so that you can provide your own implementation if needed. + */ +public class EurekaInstanceDiscovery implements InstanceDiscovery { + + private static final Logger logger = LoggerFactory.getLogger(EurekaInstanceDiscovery.class); + + // Property the controls the list of applications that are enabled in Eureka + private static final DynamicStringProperty ApplicationList = DynamicPropertyFactory.getInstance().getStringProperty("turbine.appConfig", ""); + + public EurekaInstanceDiscovery() { + // Eureka client should already be configured by spring-platform-netflix-core + // initialize eureka client. make sure eureka properties are properly configured in config.properties + //DiscoveryManager.getInstance().initComponent(new MyDataCenterInstanceConfig(), new DefaultEurekaClientConfig()); + } + + /** + * Method that queries Eureka service for a list of configured application names + * @return Collection + */ + @Override + public Collection getInstanceList() throws Exception { + + List instances = new ArrayList<>(); + + List appNames = parseApps(); + if (appNames == null || appNames.size() == 0) { + logger.info("No apps configured, returning an empty instance list"); + return instances; + } + + logger.info("Fetching instance list for apps: " + appNames); + + for (String appName : appNames) { + try { + instances.addAll(getInstancesForApp(appName)); + } catch (Exception e) { + logger.error("Failed to fetch instances for app: " + appName + ", retrying once more", e); + try { + instances.addAll(getInstancesForApp(appName)); + } catch (Exception e1) { + logger.error("Failed again to fetch instances for app: " + appName + ", giving up", e); + } + } + } + return instances; + } + + /** + * Private helper that fetches the Instances for each application. + * @param appName + * @return List + * @throws Exception + */ + private List getInstancesForApp(String appName) throws Exception { + + List instances = new ArrayList<>(); + + logger.info("Fetching instances for app: {}", appName); + Application app = DiscoveryManager.getInstance().getDiscoveryClient().getApplication(appName); + if (app == null) { + logger.warn("Eureka returned null for app: {}", appName); + } + List instancesForApp = app.getInstances(); + + if (instancesForApp != null) { + logger.info("Received instance list for app: {} = {}", appName, instancesForApp.size()); + for (InstanceInfo iInfo : instancesForApp) { + Instance instance = marshallInstanceInfo(iInfo); + if (instance != null) { + instances.add(instance); + } + } + } + + return instances; + } + + /** + * Private helper that marshals the information from each instance into something that Turbine can understand. + * Override this method for your own implementation for parsing Eureka info. + * + * @param iInfo + * @return Instance + */ + protected Instance marshallInstanceInfo(InstanceInfo iInfo) { + + String hostname = iInfo.getHostName(); + String cluster = getClusterName(iInfo); + Boolean status = parseInstanceStatus(iInfo.getStatus()); + + if (hostname != null && cluster != null && status != null) { + Instance instance = new Instance(hostname, cluster, status); + Map metadata = iInfo.getMetadata(); + if (metadata != null) { + instance.getAttributes().putAll(metadata); + } + + String asgName = iInfo.getASGName(); + if (asgName != null) { + instance.getAttributes().put("asg", asgName); + } + instance.getAttributes().put("port", String.valueOf(iInfo.getPort())); + + DataCenterInfo dcInfo = iInfo.getDataCenterInfo(); + if (dcInfo != null && dcInfo.getName().equals(DataCenterInfo.Name.Amazon)) { + AmazonInfo amznInfo = (AmazonInfo) dcInfo; + instance.getAttributes().putAll(amznInfo.getMetadata()); + } + + return instance; + } else { + return null; + } + } + + /** + * Helper that returns whether the instance is Up of Down + * @param status + * @return + */ + protected Boolean parseInstanceStatus(InstanceStatus status) { + + if (status != null) { + if (status == InstanceStatus.UP) { + return Boolean.TRUE; + } else { + return Boolean.FALSE; + } + } else { + return null; + } + } + + + /** + * Helper that fetches the cluster name. Cluster is a Turbine concept and not a Eureka concept. + * By default we choose the amazon asg name as the cluster. A custom implementation can be plugged in by overriding this method. + * + * @param iInfo + * @return + */ + protected String getClusterName(InstanceInfo iInfo) { + return iInfo.getASGName(); + } + + /** + * TODO: move to ConfigurationProperties + * Private helper that parses the list of application names. + * + * @return List + */ + private List parseApps() { + + String appList = ApplicationList.get(); + if (appList == null) { + return null; + } + + appList = appList.trim(); + if (appList.length() == 0) { + return null; + } + + String[] parts = appList.split(","); + if (parts != null && parts.length > 0) { + return Arrays.asList(parts); + } + + return null; + } +} \ No newline at end of file diff --git a/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/SpringAggregatorFactory.java b/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/SpringAggregatorFactory.java new file mode 100644 index 00000000..19a86d3e --- /dev/null +++ b/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/SpringAggregatorFactory.java @@ -0,0 +1,121 @@ +package io.spring.platform.netflix.turbine; + +import com.netflix.config.DynamicPropertyFactory; +import com.netflix.config.DynamicStringProperty; +import com.netflix.turbine.data.AggDataFromCluster; +import com.netflix.turbine.discovery.Instance; +import com.netflix.turbine.handler.PerformanceCriteria; +import com.netflix.turbine.handler.TurbineDataHandler; +import com.netflix.turbine.monitor.TurbineDataMonitor; +import com.netflix.turbine.monitor.cluster.AggregateClusterMonitor; +import com.netflix.turbine.monitor.cluster.ClusterMonitor; +import com.netflix.turbine.monitor.cluster.ClusterMonitorFactory; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; + +import static com.netflix.turbine.monitor.cluster.AggregateClusterMonitor.AggregatorClusterMonitorConsole; + +/** + * Created by sgibb on 7/14/14. + */ +public class SpringAggregatorFactory implements ClusterMonitorFactory { + private static final Logger logger = LoggerFactory.getLogger(SpringAggregatorFactory.class); + + private static final DynamicStringProperty aggClusters = DynamicPropertyFactory.getInstance().getStringProperty("turbine.aggregator.clusterConfig", null); + + /** + * @return {@link ClusterMonitor}<{@link AggDataFromCluster}> + */ + @Override + public ClusterMonitor getClusterMonitor(String name) { + TurbineDataMonitor clusterMonitor = AggregateClusterMonitor.AggregatorClusterMonitorConsole.findMonitor(name + "_agg"); + return (ClusterMonitor) clusterMonitor; + } + + public static TurbineDataMonitor findOrRegisterAggregateMonitor(String clusterName) { + + TurbineDataMonitor clusterMonitor = AggregatorClusterMonitorConsole.findMonitor(clusterName + "_agg"); + + if (clusterMonitor == null) { + logger.info("Could not find monitors: " + AggregatorClusterMonitorConsole.toString()); + clusterMonitor = new SpringClusterMonitor(clusterName + "_agg", clusterName); + clusterMonitor = AggregatorClusterMonitorConsole.findOrRegisterMonitor(clusterMonitor); + } + + return clusterMonitor; + } + + @Override + public void initClusterMonitors() { + for(String clusterName : getClusterNames()) { + ClusterMonitor clusterMonitor = (ClusterMonitor) findOrRegisterAggregateMonitor(clusterName); + clusterMonitor.registerListenertoClusterMonitor(StaticListener); + try { + clusterMonitor.startMonitor(); + } catch (Exception e) { + logger.warn("Could not init cluster monitor for: " + clusterName); + clusterMonitor.stopMonitor(); + clusterMonitor.getDispatcher().stopDispatcher(); + } + } + } + + private List getClusterNames() { + + List clusters = new ArrayList(); + String clusterNames = aggClusters.get(); + if (clusterNames == null || clusterNames.trim().length() == 0) { + clusters.add("default"); + } else { + String[] parts = aggClusters.get().split(","); + for (String s : parts) { + clusters.add(s); + } + } + return clusters; + } + + private TurbineDataHandler StaticListener = new TurbineDataHandler() { + + @Override + public String getName() { + return "StaticListener_For_Aggregator"; + } + + @Override + public void handleData(Collection stats) { + } + + @Override + public void handleHostLost(Instance host) { + } + + @Override + public PerformanceCriteria getCriteria() { + return NonCriticalCriteria; + } + + }; + + private PerformanceCriteria NonCriticalCriteria = new PerformanceCriteria() { + + @Override + public boolean isCritical() { + return false; + } + + @Override + public int getMaxQueueSize() { + return 0; + } + + @Override + public int numThreads() { + return 0; + } + }; +} diff --git a/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/SpringClusterMonitor.java b/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/SpringClusterMonitor.java new file mode 100644 index 00000000..a6882ef5 --- /dev/null +++ b/spring-platform-netflix-turbine/src/main/java/io/spring/platform/netflix/turbine/SpringClusterMonitor.java @@ -0,0 +1,79 @@ +package io.spring.platform.netflix.turbine; + +import com.netflix.config.DynamicBooleanProperty; +import com.netflix.config.DynamicPropertyFactory; +import com.netflix.config.DynamicStringProperty; +import com.netflix.turbine.data.DataFromSingleInstance; +import com.netflix.turbine.discovery.Instance; +import com.netflix.turbine.handler.PerformanceCriteria; +import com.netflix.turbine.monitor.MonitorConsole; +import com.netflix.turbine.monitor.cluster.AggregateClusterMonitor; +import com.netflix.turbine.monitor.cluster.ObservationCriteria; +import com.netflix.turbine.monitor.instance.InstanceUrlClosure; + +/** + * Created by sgibb on 7/14/14. + */ +public class SpringClusterMonitor extends AggregateClusterMonitor { + + public SpringClusterMonitor(String name, String clusterName) { + super(name, + new ObservationCriteria.ClusterBasedObservationCriteria(clusterName), + new PerformanceCriteria.AggClusterPerformanceCriteria(clusterName), + new MonitorConsole(), + InstanceMonitorDispatcher, + SpringClusterMonitor.ClusterConfigBasedUrlClosure); + } + + /** + * TODO: make this a template of some kind (secure, management port, etc...) + * Helper class that decides how to connect to a server based on injected config. + * Note that the cluster name must be provided here since one can have different configs for different clusters + */ + public static InstanceUrlClosure ClusterConfigBasedUrlClosure = new InstanceUrlClosure() { + + private final DynamicStringProperty defaultUrlClosureConfig = DynamicPropertyFactory.getInstance().getStringProperty("turbine.instanceUrlSuffix", null); + private final DynamicBooleanProperty instanceInsertPort = DynamicPropertyFactory.getInstance().getBooleanProperty("turbine.instanceInsertPort", false); + @Override + public String getUrlPath(Instance host) { + + if (host.getCluster() == null) { + throw new RuntimeException("Host must have cluster name in order to use ClusterConfigBasedUrlClosure"); + } + + String key = "turbine.instanceUrlSuffix." + host.getCluster(); + DynamicStringProperty urlClosureConfig = DynamicPropertyFactory.getInstance().getStringProperty(key, null); + + String url = urlClosureConfig.get(); + if (url == null) { + url = defaultUrlClosureConfig.get(); + } + + if (url == null) { + throw new RuntimeException("Config property: " + urlClosureConfig.getName() + " or " + + defaultUrlClosureConfig.getName() + " must be set"); + } + + String insertPortKey = "turbine.instanceInsertPort." + host.getCluster(); + DynamicStringProperty insertPortProp = DynamicPropertyFactory.getInstance().getStringProperty(insertPortKey, null); + boolean insertPort; + if (insertPortProp.get() == null) { + insertPort = instanceInsertPort.get(); + } else { + insertPort = Boolean.parseBoolean(insertPortProp.get()); + } + + if (insertPort) { + if (url.startsWith("/")) { + url = url.substring(1); + } + if (!host.getAttributes().containsKey("port")) { + throw new RuntimeException("Configured to use port, but port is not in host attributes"); + } + return String.format("http://%s:%s/%s", host.getHostname(), host.getAttributes().get("port"), url); + } + + return "http://" + host.getHostname() + url; + } + }; +}