diff --git a/labx-13/labx-13-sc-sleuth-db-elasticsearch/pom.xml b/labx-13/labx-13-sc-sleuth-db-elasticsearch/pom.xml new file mode 100644 index 00000000..68600057 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-db-elasticsearch/pom.xml @@ -0,0 +1,78 @@ + + + + labx-13 + cn.iocoder.springboot.labs + 1.0-SNAPSHOT + + 4.0.0 + + labx-13-sc-sleuth-db-elasticsearch + + + 1.8 + 1.8 + 2.2.4.RELEASE + Hoxton.SR1 + + + + + + + org.springframework.boot + spring-boot-starter-parent + ${spring.boot.version} + pom + import + + + org.springframework.cloud + spring-cloud-dependencies + ${spring.cloud.version} + pom + import + + + + + + + + org.springframework.boot + spring-boot-starter-web + + + + + org.springframework.boot + spring-boot-starter-data-elasticsearch + + + + + org.springframework.cloud + spring-cloud-starter-zipkin + + + + + io.opentracing.brave + brave-opentracing + 0.35.0 + + + + + io.opentracing.contrib + opentracing-elasticsearch6-client + 0.1.6 + + + + diff --git a/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/UserServiceApplication.java b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/UserServiceApplication.java new file mode 100644 index 00000000..5386ea90 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/UserServiceApplication.java @@ -0,0 +1,13 @@ +package cn.iocoder.springcloud.labx13.springmvcdemo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class UserServiceApplication { + + public static void main(String[] args) { + SpringApplication.run(UserServiceApplication.class, args); + } + +} diff --git a/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/config/SleuthConfiguration.java b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/config/SleuthConfiguration.java new file mode 100644 index 00000000..d5974d81 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/config/SleuthConfiguration.java @@ -0,0 +1,35 @@ +package cn.iocoder.springcloud.labx13.springmvcdemo.config; + +import cn.iocoder.springcloud.labx13.springmvcdemo.spring.TracingTransportClientFactoryBean; +import io.opentracing.Tracer; +import org.elasticsearch.client.transport.TransportClient; +import org.springframework.boot.autoconfigure.data.elasticsearch.ElasticsearchProperties; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +import java.util.Properties; + +@Configuration +public class SleuthConfiguration { + + // ==================== Elasticsearch 相关 ==================== + + @Bean + public TransportClient elasticsearchClient(Tracer tracer, ElasticsearchProperties elasticsearchProperties) throws Exception { + // 创建 TracingTransportClientFactoryBean 对象 + TracingTransportClientFactoryBean factory = new TracingTransportClientFactoryBean(tracer); + // 设置其属性 + factory.setClusterNodes(elasticsearchProperties.getClusterNodes()); + factory.setProperties(this.createElasticsearch(elasticsearchProperties)); + // 创建 TransportClient 对象,并返回 + factory.afterPropertiesSet(); + return factory.getObject(); + } + + private Properties createElasticsearch(ElasticsearchProperties elasticsearchProperties) { + Properties properties = new Properties(); + properties.put("cluster.name", elasticsearchProperties.getClusterName()); + properties.putAll(elasticsearchProperties.getProperties()); + return properties; + } +} diff --git a/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/controller/UserController.java b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/controller/UserController.java new file mode 100644 index 00000000..8ba647a3 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/controller/UserController.java @@ -0,0 +1,28 @@ +package cn.iocoder.springcloud.labx13.springmvcdemo.controller; + +import cn.iocoder.springcloud.labx13.springmvcdemo.dataobject.ESUserDO; +import cn.iocoder.springcloud.labx13.springmvcdemo.repository.ESUserRepository; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RequestParam; +import org.springframework.web.bind.annotation.RestController; + +@RestController +@RequestMapping("/user") +public class UserController { + + @Autowired + private ESUserRepository userRepository; + + @GetMapping("/get") + public String get(@RequestParam("id") Integer id) { + this.findById(1); + return "success"; + } + + public ESUserDO findById(Integer id) { + return userRepository.findById(id).orElse(null); + } + +} diff --git a/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/dataobject/ESUserDO.java b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/dataobject/ESUserDO.java new file mode 100644 index 00000000..df9ab748 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/dataobject/ESUserDO.java @@ -0,0 +1,41 @@ +package cn.iocoder.springcloud.labx13.springmvcdemo.dataobject; + +import org.springframework.data.annotation.Id; +import org.springframework.data.elasticsearch.annotations.Document; + +import java.util.Date; + +@Document(indexName = "user", // 索引名 + type = "user", // 类型。未来的版本即将废弃 + shards = 1, // 默认索引分区数 + replicas = 0, // 每个分区的备份数 + refreshInterval = "-1" // 刷新间隔 +) +public class ESUserDO { + + @Id + private Integer id; + /** + * 账号 + */ + private String username; + /** + * 密码 + */ + private String password; + /** + * 创建时间 + */ + private Date createTime; + + @Override + public String toString() { + return "UserDO{" + + "id=" + id + + ", username='" + username + '\'' + + ", password='" + password + '\'' + + ", createTime=" + createTime + + '}'; + } + +} diff --git a/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/repository/ESUserRepository.java b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/repository/ESUserRepository.java new file mode 100644 index 00000000..005a24b8 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/repository/ESUserRepository.java @@ -0,0 +1,8 @@ +package cn.iocoder.springcloud.labx13.springmvcdemo.repository; + +import cn.iocoder.springcloud.labx13.springmvcdemo.dataobject.ESUserDO; +import org.springframework.data.elasticsearch.repository.ElasticsearchRepository; + +public interface ESUserRepository extends ElasticsearchRepository { + +} diff --git a/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/spring/ClusterNodes.java b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/spring/ClusterNodes.java new file mode 100644 index 00000000..ecc9d831 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/spring/ClusterNodes.java @@ -0,0 +1,81 @@ +package cn.iocoder.springcloud.labx13.springmvcdemo.spring; + +import org.elasticsearch.common.transport.TransportAddress; +import org.springframework.data.util.Streamable; +import org.springframework.util.Assert; +import org.springframework.util.StringUtils; + +import java.net.InetAddress; +import java.net.UnknownHostException; +import java.util.Arrays; +import java.util.Iterator; +import java.util.List; +import java.util.stream.Collectors; + +class ClusterNodes implements Streamable { + + public static ClusterNodes DEFAULT = ClusterNodes.of("127.0.0.1:9300"); + + private static final String COLON = ":"; + private static final String COMMA = ","; + + private final List clusterNodes; + + /** + * Creates a new {@link ClusterNodes} by parsing the given source. + * + * @param source must not be {@literal null} or empty. + */ + private ClusterNodes(String source) { + + Assert.hasText(source, "Cluster nodes source must not be null or empty!"); + + String[] nodes = StringUtils.delimitedListToStringArray(source, COMMA); + + this.clusterNodes = Arrays.stream(nodes).map(node -> { + + String[] segments = StringUtils.delimitedListToStringArray(node, COLON); + + Assert.isTrue(segments.length == 2, + () -> String.format("Invalid cluster node %s in %s! Must be in the format host:port!", node, source)); + + String host = segments[0].trim(); + String port = segments[1].trim(); + + Assert.hasText(host, () -> String.format("No host name given cluster node %s!", node)); + Assert.hasText(port, () -> String.format("No port given in cluster node %s!", node)); + + return new TransportAddress(toInetAddress(host), Integer.valueOf(port)); + + }).collect(Collectors.toList()); + } + + /** + * Creates a new {@link ClusterNodes} by parsing the given source. The expected format is a comma separated list of + * host-port-combinations separated by a colon: {@code host:port,host:port,…}. + * + * @param source must not be {@literal null} or empty. + * @return + */ + public static ClusterNodes of(String source) { + return new ClusterNodes(source); + } + + /* + * (non-Javadoc) + * @see java.lang.Iterable#iterator() + */ + @Override + public Iterator iterator() { + return clusterNodes.iterator(); + } + + private static InetAddress toInetAddress(String host) { + + try { + return InetAddress.getByName(host); + } catch (UnknownHostException o_O) { + throw new IllegalArgumentException(o_O); + } + } +} diff --git a/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/spring/TracingTransportClientFactoryBean.java b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/spring/TracingTransportClientFactoryBean.java new file mode 100644 index 00000000..27646b52 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/java/cn/iocoder/springcloud/labx13/springmvcdemo/spring/TracingTransportClientFactoryBean.java @@ -0,0 +1,138 @@ +package cn.iocoder.springcloud.labx13.springmvcdemo.spring; + +import io.opentracing.Tracer; +import io.opentracing.contrib.elasticsearch6.TracingPreBuiltTransportClient; +import org.elasticsearch.client.transport.TransportClient; +import org.elasticsearch.common.settings.Settings; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.DisposableBean; +import org.springframework.beans.factory.FactoryBean; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.data.elasticsearch.client.TransportClientFactoryBean; + +import java.util.Properties; + +/** + * 参考 {@link TransportClientFactoryBean} 来实现。 + */ +public class TracingTransportClientFactoryBean implements FactoryBean, InitializingBean, DisposableBean { + + private static final Logger logger = LoggerFactory.getLogger(TransportClientFactoryBean.class); + private ClusterNodes clusterNodes = ClusterNodes.of("127.0.0.1:9300"); + private String clusterName = "elasticsearch"; + private Boolean clientTransportSniff = true; + private Boolean clientIgnoreClusterName = Boolean.FALSE; + private String clientPingTimeout = "5s"; + private String clientNodesSamplerInterval = "5s"; + private TransportClient client; + private Properties properties; + + private Tracer tracer; + + public TracingTransportClientFactoryBean(Tracer tracer) { + this.tracer = tracer; + } + + @Override + public void destroy() throws Exception { + try { + logger.info("Closing elasticSearch client"); + if (client != null) { + client.close(); + } + } catch (final Exception e) { + logger.error("Error closing ElasticSearch client: ", e); + } + } + + @Override + public TransportClient getObject() throws Exception { + return client; + } + + @Override + public Class getObjectType() { + return TransportClient.class; + } + + @Override + public boolean isSingleton() { + return true; + } + + @Override + public void afterPropertiesSet() throws Exception { + buildClient(); + } + + protected void buildClient() throws Exception { + // 创建可追踪的 TracingPreBuiltTransportClient + client = new TracingPreBuiltTransportClient(tracer, settings()); + + clusterNodes.stream() // + .peek(it -> logger.info("Adding transport node : " + it.toString())) // + .forEach(client::addTransportAddress); + + client.connectedNodes(); + } + + private Settings settings() { + if (properties != null) { + Settings.Builder builder = Settings.builder(); + + properties.forEach((key, value) -> { + builder.put(key.toString(), value.toString()); + }); + + return builder.build(); + } + return Settings.builder() + .put("cluster.name", clusterName) + .put("client.transport.sniff", clientTransportSniff) + .put("client.transport.ignore_cluster_name", clientIgnoreClusterName) + .put("client.transport.ping_timeout", clientPingTimeout) + .put("client.transport.nodes_sampler_interval", clientNodesSamplerInterval) + .build(); + } + + public void setClusterNodes(String clusterNodes) { + this.clusterNodes = ClusterNodes.of(clusterNodes); + } + + public void setClusterName(String clusterName) { + this.clusterName = clusterName; + } + + public void setClientTransportSniff(Boolean clientTransportSniff) { + this.clientTransportSniff = clientTransportSniff; + } + + public String getClientNodesSamplerInterval() { + return clientNodesSamplerInterval; + } + + public void setClientNodesSamplerInterval(String clientNodesSamplerInterval) { + this.clientNodesSamplerInterval = clientNodesSamplerInterval; + } + + public String getClientPingTimeout() { + return clientPingTimeout; + } + + public void setClientPingTimeout(String clientPingTimeout) { + this.clientPingTimeout = clientPingTimeout; + } + + public Boolean getClientIgnoreClusterName() { + return clientIgnoreClusterName; + } + + public void setClientIgnoreClusterName(Boolean clientIgnoreClusterName) { + this.clientIgnoreClusterName = clientIgnoreClusterName; + } + + public void setProperties(Properties properties) { + this.properties = properties; + } +} diff --git a/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/resources/application.yml b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/resources/application.yml new file mode 100644 index 00000000..3c812e67 --- /dev/null +++ b/labx-13/labx-13-sc-sleuth-db-elasticsearch/src/main/resources/application.yml @@ -0,0 +1,19 @@ +spring: + application: + name: user-service # 服务名 + + # Zipkin 配置项,对应 ZipkinProperties 类 + zipkin: + base-url: http://127.0.0.1:9411 # Zipkin 服务的地址 + + # Spring Cloud Sleuth 配置项 + sleuth: + # Spring Cloud Sleuth 针对 Web 组件的配置项,例如说 SpringMVC + web: + enabled: true # 是否开启,默认为 true + + data: + # Elasticsearch 配置项 + elasticsearch: + cluster-name: elasticsearch # 集群名 + cluster-nodes: 127.0.0.1:9300 # 集群节点 diff --git a/labx-13/pom.xml b/labx-13/pom.xml index 3991fe0d..2a473bfd 100644 --- a/labx-13/pom.xml +++ b/labx-13/pom.xml @@ -22,6 +22,7 @@ labx-13-sc-sleuth-db-mysql labx-13-sc-sleuth-db-redis labx-13-sc-sleuth-db-mongodb + labx-13-sc-sleuth-db-elasticsearch