spring-35-36-spring-cloud

This commit is contained in:
petrelevich
2025-12-02 21:38:02 +03:00
parent 16920a4326
commit f061675f1e
41 changed files with 1170 additions and 0 deletions
+121
View File
@@ -0,0 +1,121 @@
import io.spring.gradle.dependencymanagement.dsl.DependencyManagementExtension
import name.remal.gradle_plugins.sonarlint.SonarLintExtension
import org.springframework.boot.gradle.plugin.SpringBootPlugin.BOM_COORDINATES
import org.gradle.plugins.ide.idea.model.IdeaLanguageLevel
import org.springframework.boot.gradle.plugin.SpringBootPlugin
plugins {
idea
id("fr.brouillard.oss.gradle.jgitver")
id("io.spring.dependency-management")
id("org.springframework.boot") apply false
id("name.remal.sonarlint") apply false
id("com.diffplug.spotless") apply false
}
allprojects {
group = "ru.demo"
repositories {
mavenLocal()
mavenCentral()
}
val springCloudVersion: String by project
val logbackEncoder: String by project
apply(plugin = "io.spring.dependency-management")
dependencyManagement {
dependencies {
imports {
mavenBom(BOM_COORDINATES)
mavenBom("org.springframework.cloud:spring-cloud-dependencies:$springCloudVersion")
}
dependency("net.logstash.logback:logstash-logback-encoder:$logbackEncoder")
}
}
configurations.all {
resolutionStrategy {
failOnVersionConflict()
force("org.glassfish.hk2.external:aopalliance-repackaged:3.1.1")
force("org.glassfish.hk2:hk2-utils:3.1.1")
force("commons-logging:commons-logging:1.3.1")
force("com.fasterxml.woodstox:woodstox-core:6.6.2")
force("org.glassfish.hk2:hk2-api:3.1.1")
force("org.apache.httpcomponents:httpclient:4.5.14")
force("org.sonarsource.analyzer-commons:sonar-analyzer-commons:2.11.0.2861")
force("org.sonarsource.analyzer-commons:sonar-xml-parsing:2.11.0.2861")
force("org.sonarsource.sslr:sslr-core:1.24.0.633")
force("org.sonarsource.analyzer-commons:sonar-analyzer-recognizers:2.11.0.2861")
force("commons-io:commons-io:2.15.1")
force("com.google.guava:guava:32.1.3-jre")
force("com.google.code.findbugs:jsr305:3.0.2")
force("org.codehaus.woodstox:stax2-api:4.2.2")
force("io.opentelemetry:opentelemetry-api-incubator:1.38.0-alpha")
}
}
}
subprojects {
plugins.apply(SpringBootPlugin::class.java)
plugins.apply(JavaPlugin::class.java)
extensions.configure<JavaPluginExtension> {
sourceCompatibility = JavaVersion.VERSION_21
targetCompatibility = JavaVersion.VERSION_21
}
tasks.withType<JavaCompile> {
options.encoding = "UTF-8"
options.compilerArgs.addAll(listOf("-Xlint:all,-serial,-processing"))
dependsOn("spotlessApply")
}
apply<name.remal.gradle_plugins.sonarlint.SonarLintPlugin>()
configure<SonarLintExtension> {
nodeJs {
detectNodeJs = false
logNodeJsNotFound = false
}
}
apply<com.diffplug.gradle.spotless.SpotlessPlugin>()
configure<com.diffplug.gradle.spotless.SpotlessExtension> {
java {
palantirJavaFormat("2.38.0")
}
}
plugins.apply(fr.brouillard.oss.gradle.plugins.JGitverPlugin::class.java)
extensions.configure<fr.brouillard.oss.gradle.plugins.JGitverPluginExtension> {
strategy("PATTERN")
nonQualifierBranches("main,master")
tagVersionPattern("\${v}\${<meta.DIRTY_TEXT}")
versionPattern(
"\${v}\${<meta.COMMIT_DISTANCE}\${<meta.GIT_SHA1_8}" +
"\${<meta.QUALIFIED_BRANCH_NAME}\${<meta.DIRTY_TEXT}-SNAPSHOT"
)
}
tasks.withType<Test> {
useJUnitPlatform()
testLogging.showExceptions = true
reports {
junitXml.required.set(true)
html.required.set(true)
}
}
}
tasks {
val managedVersions by registering {
doLast {
project.extensions.getByType<DependencyManagementExtension>()
.managedVersions
.toSortedMap()
.map { "${it.key}:${it.value}" }
.forEach(::println)
}
}
}
@@ -0,0 +1,38 @@
plugins {
id("com.google.cloud.tools.jib")
}
dependencies {
implementation(project(":kafka-log-appender"))
implementation("net.logstash.logback:logstash-logback-encoder")
implementation("org.springframework.boot:spring-boot-starter-actuator")
implementation("io.micrometer:micrometer-registry-prometheus")
implementation("org.springframework.cloud:spring-cloud-starter-config")
implementation ("org.springframework.cloud:spring-cloud-starter-netflix-eureka-server")
implementation("io.micrometer:micrometer-tracing-bridge-otel") // bridges the Micrometer Observation API to OpenTelemetry.
implementation("io.opentelemetry:opentelemetry-exporter-zipkin") // reports traces to Zipkin.
}
jib {
container {
creationTime.set("USE_CURRENT_TIMESTAMP")
}
from {
image = "bellsoft/liberica-openjdk-alpine-musl:21.0.1"
}
to {
image = "localrun/eureka-server"
tags = setOf(project.version.toString())
}
}
tasks {
build {
dependsOn(spotlessApply)
dependsOn(jibBuildTar)
}
}
@@ -0,0 +1,18 @@
#!/bin/bash
../gradlew :eureka-server:build
docker load --input build/jib-image.tar
docker stop eureka-server
docker run --rm -d --name eureka-server \
--memory=512m \
--cpus 1 \
--network="host" \
-v $HOME/.ssh:/root/.ssh \
-e JAVA_TOOL_OPTIONS="-XX:InitialRAMPercentage=80 -XX:MaxRAMPercentage=80" \
localrun/eureka-server:latest
+15
View File
@@ -0,0 +1,15 @@
# -------Gradle--------
org.gradle.jvmargs=-Xmx4g
org.gradle.daemon=true
org.gradle.parallel=true
# -------Plugins---------
jgitver=0.10.0-rc03
dependencyManagement=1.1.5
springframeworkBoot=3.3.1
springCloudVersion=2023.0.3
sonarlint=4.2.4
spotless=6.25.0
jib=3.4.5
# -------Versions--------
logbackEncoder=8.0
@@ -0,0 +1,10 @@
dependencies {
implementation("ch.qos.logback:logback-classic")
implementation ("org.apache.kafka:kafka-clients")
}
tasks {
bootJar {
enabled = false
}
}
@@ -0,0 +1,8 @@
package ru.appender.kafka;
public class AppenderException extends RuntimeException {
public AppenderException(String message) {
super(message);
}
}
@@ -0,0 +1,5 @@
package ru.appender.kafka;
import java.util.function.Consumer;
public interface ErrorMsgConsumer extends Consumer<String> {}
@@ -0,0 +1,84 @@
package ru.appender.kafka;
import ch.qos.logback.classic.spi.LoggingEvent;
import ch.qos.logback.core.UnsynchronizedAppenderBase;
import ch.qos.logback.core.encoder.Encoder;
import java.util.Queue;
import java.util.concurrent.ArrayBlockingQueue;
public class LogAppender extends UnsynchronizedAppenderBase<LoggingEvent> {
private static final String MESSAGE_TEMPLATE = "[Kafka appender] %s";
private String bootstrapServers;
private String topicName;
private final Queue<LoggingEvent> eventsQueue = new ArrayBlockingQueue<>(1000);
private Thread senderThread;
private Encoder<LoggingEvent> encoder;
private final ErrorMsgConsumer errorMsgConsumer = error -> addError(String.format(MESSAGE_TEMPLATE, error));
public void setEncoder(Encoder<LoggingEvent> encoder) {
this.encoder = encoder;
}
public void setBootstrapServers(String bootstrapServers) {
this.bootstrapServers = bootstrapServers;
addInfo(String.format(MESSAGE_TEMPLATE, "set bootstrapServers:" + bootstrapServers));
}
public void setTopicName(String topicName) {
this.topicName = topicName;
addInfo(String.format(MESSAGE_TEMPLATE, "set topicName:" + topicName));
}
@Override
public void start() {
if (bootstrapServers == null) {
addError(String.format(MESSAGE_TEMPLATE, "bootstrapServers is null"));
return;
}
if (topicName == null) {
addError(String.format(MESSAGE_TEMPLATE, "topicName is null"));
return;
}
var logProducer = new LogProducer(bootstrapServers, topicName);
senderThread = Thread.ofVirtual().name("senderThread").start(() -> sendMessages(logProducer));
super.start();
addInfo(String.format(MESSAGE_TEMPLATE, "started"));
}
@Override
public void stop() {
super.stop();
if (senderThread != null) {
senderThread.interrupt();
}
addInfo(String.format(MESSAGE_TEMPLATE, "stopped"));
}
@Override
public void append(LoggingEvent eventObject) {
var result = eventsQueue.offer(eventObject);
if (!result) {
addWarn(String.format(MESSAGE_TEMPLATE, "eventsQueue is full"));
}
}
private void sendMessages(LogProducer logProducer) {
while (!Thread.currentThread().isInterrupted()) {
var event = eventsQueue.poll();
if (event != null) {
try {
var messageAsText = new String(encoder.encode(event));
logProducer.send(messageAsText, errorMsgConsumer);
} catch (Exception ex) {
addError(String.format(MESSAGE_TEMPLATE, ex.getMessage()));
}
}
}
}
}
@@ -0,0 +1,56 @@
package ru.appender.kafka;
import static org.apache.kafka.clients.CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG;
import static org.apache.kafka.clients.CommonClientConfigs.CLIENT_ID_CONFIG;
import static org.apache.kafka.clients.CommonClientConfigs.RETRIES_CONFIG;
import static org.apache.kafka.clients.producer.ProducerConfig.ACKS_CONFIG;
import static org.apache.kafka.clients.producer.ProducerConfig.BATCH_SIZE_CONFIG;
import static org.apache.kafka.clients.producer.ProducerConfig.BUFFER_MEMORY_CONFIG;
import static org.apache.kafka.clients.producer.ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG;
import static org.apache.kafka.clients.producer.ProducerConfig.LINGER_MS_CONFIG;
import static org.apache.kafka.clients.producer.ProducerConfig.MAX_BLOCK_MS_CONFIG;
import static org.apache.kafka.clients.producer.ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG;
import java.util.Properties;
import java.util.function.Consumer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;
public class LogProducer {
private final KafkaProducer<String, String> kafkaProducer;
private final String topicName;
private long lastSendKey = System.currentTimeMillis();
public LogProducer(String bootstrapServers, String topicName) {
this.topicName = topicName;
Properties props = new Properties();
props.put(CLIENT_ID_CONFIG, "myKafkaProducer");
props.put(BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ACKS_CONFIG, "1");
props.put(RETRIES_CONFIG, 1);
props.put(BATCH_SIZE_CONFIG, 16384);
props.put(LINGER_MS_CONFIG, 10);
props.put(BUFFER_MEMORY_CONFIG, 33_554_432); // bytes
props.put(MAX_BLOCK_MS_CONFIG, 1_000); // ms
props.put(KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
kafkaProducer = new KafkaProducer<>(props);
Runtime.getRuntime().addShutdownHook(new Thread(this::close));
}
public void send(String value, Consumer<String> errorCallback) {
var key = lastSendKey++;
kafkaProducer.send(new ProducerRecord<>(topicName, String.valueOf(key), value), (metadata, exception) -> {
if (exception != null) {
errorCallback.accept(String.format(exception.getMessage()));
}
});
}
public void close() {
kafkaProducer.close();
}
}
@@ -0,0 +1,14 @@
dependencies {
implementation(project(":kafka-log-appender"))
implementation("net.logstash.logback:logstash-logback-encoder")
implementation ("org.springframework.boot:spring-boot-starter-web")
implementation("org.springframework.boot:spring-boot-starter-actuator")
implementation("io.micrometer:micrometer-registry-prometheus")
implementation("org.springframework.cloud:spring-cloud-starter-config")
implementation("org.springframework.cloud:spring-cloud-starter-netflix-eureka-client")
implementation("io.micrometer:micrometer-tracing-bridge-otel") // bridges the Micrometer Observation API to OpenTelemetry.
implementation("io.opentelemetry:opentelemetry-exporter-zipkin") // reports traces to Zipkin.
}
@@ -0,0 +1,11 @@
package ru.demo;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.builder.SpringApplicationBuilder;
@SpringBootApplication
public class ServiceClientInfo {
public static void main(String[] args) {
new SpringApplicationBuilder().sources(ServiceClientInfo.class).run(args);
}
}
@@ -0,0 +1,19 @@
package ru.demo.config;
import org.springframework.boot.web.servlet.FilterRegistrationBean;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.filter.OncePerRequestFilter;
import ru.demo.filter.MdcFilter;
@Configuration
public class ApplConf {
@Bean
public FilterRegistrationBean<OncePerRequestFilter> mdcFilterRegistrationBean() {
var registrationBean = new FilterRegistrationBean<OncePerRequestFilter>();
registrationBean.setFilter(new MdcFilter());
registrationBean.setOrder(1);
return registrationBean;
}
}
@@ -0,0 +1,30 @@
package ru.demo.controller;
import java.time.Duration;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import ru.demo.model.ClientData;
@RestController
public class ClientInfoController {
private static final Logger logger = LoggerFactory.getLogger(ClientInfoController.class);
// curl -v http://localhost:8083/additional-info?name="testClient"
// docker run -d --rm -p 9411:9411 openzipkin/zipkin
// http://localhost:9411
@GetMapping(value = "/additional-info")
public ClientData info(@RequestParam(name = "name") String name) throws InterruptedException {
logger.info("request. name:{}", name);
doJob();
return new ClientData(String.format("additional ClientInfo name:%s", name));
}
private void doJob() throws InterruptedException {
Thread.sleep(Duration.ofMillis(100));
}
}
@@ -0,0 +1,39 @@
package ru.demo.filter;
import jakarta.servlet.FilterChain;
import jakarta.servlet.ServletException;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.servlet.http.HttpServletResponse;
import java.io.IOException;
import java.util.ArrayList;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.MDC;
import org.springframework.web.filter.OncePerRequestFilter;
public class MdcFilter extends OncePerRequestFilter {
private final Logger log = LoggerFactory.getLogger(MdcFilter.class);
private static final String HEADER_X_REQUEST_ID = "X-Request-Id";
private static final String MDC_REQUEST_ID = "requestId";
@Override
protected void doFilterInternal(HttpServletRequest request, HttpServletResponse response, FilterChain filterChain)
throws ServletException, IOException {
var xRequestId = request.getHeader(HEADER_X_REQUEST_ID);
log.info("method:{}, xRequestId:{}", request.getMethod(), xRequestId);
if (xRequestId != null) {
MDC.put(MDC_REQUEST_ID, xRequestId);
}
var headerIterator = request.getHeaderNames().asIterator();
var headers = new ArrayList<String>();
while (headerIterator.hasNext()) {
headers.add(headerIterator.next());
}
log.info("request headers:{}", headers);
response.addHeader(HEADER_X_REQUEST_ID, xRequestId);
filterChain.doFilter(request, response);
MDC.remove(MDC_REQUEST_ID);
}
}
@@ -0,0 +1,3 @@
package ru.demo.model;
public record ClientData(String data) {}
@@ -0,0 +1,16 @@
spring:
application:
name: service-client-info
cloud:
config:
fail-fast: true
retry:
initial-interval: 5000
max-attempts: 10
max-interval: 5000
multiplier: 1.2
config:
import: optional:configserver:http://localhost:8888
codec:
max-in-memory-size: 10MB
@@ -0,0 +1,26 @@
<configuration scan="true" scanPeriod="10 seconds">
<jmxConfigurator />
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - [%X{requestId}] %msg%n</pattern>
</encoder>
</appender>
<appender name="kafka" class="ru.appender.kafka.LogAppender">
<bootstrapServers>localhost:9092</bootstrapServers>
<topicName>applLogs</topicName>
<encoder class="net.logstash.logback.encoder.LogstashEncoder">
<fieldNames>
<version>[ignore]</version>
</fieldNames>
<customFields>{"appname":"service-client-info"}</customFields>
<includeMdcKeyName>requestId</includeMdcKeyName>
</encoder>
</appender>
<root level="info">
<appender-ref ref="STDOUT" />
<appender-ref ref="kafka" />
</root>
</configuration>
@@ -0,0 +1,19 @@
dependencies {
implementation(project(":kafka-log-appender"))
implementation("net.logstash.logback:logstash-logback-encoder")
implementation("org.springframework.boot:spring-boot-starter-actuator")
implementation("io.micrometer:micrometer-registry-prometheus")
implementation("org.springframework.boot:spring-boot-starter-web")
implementation("org.springframework.cloud:spring-cloud-starter-config")
implementation("org.springframework.cloud:spring-cloud-starter-netflix-eureka-client")
implementation("org.springframework.cloud:spring-cloud-starter-openfeign")
implementation("io.micrometer:micrometer-tracing-bridge-otel") // bridges the Micrometer Observation API to OpenTelemetry.
implementation("io.opentelemetry:opentelemetry-exporter-zipkin") // reports traces to Zipkin.
implementation("io.github.openfeign:feign-micrometer")
implementation("org.springframework.cloud:spring-cloud-starter-circuitbreaker-resilience4j")
}
@@ -0,0 +1,12 @@
package ru.demo;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.builder.SpringApplicationBuilder;
@SpringBootApplication
public class ServiceClient {
public static void main(String[] args) {
new SpringApplicationBuilder().sources(ServiceClient.class).run(args);
}
}
@@ -0,0 +1,23 @@
package ru.demo.config;
import feign.RequestTemplate;
import feign.codec.EncodeException;
import feign.codec.Encoder;
import java.lang.reflect.Type;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class RequestEncoder implements Encoder {
private static final Logger log = LoggerFactory.getLogger(RequestEncoder.class);
private final Encoder defaultEncoder;
public RequestEncoder(Encoder defaultEncoder) {
this.defaultEncoder = defaultEncoder;
}
@Override
public void encode(Object object, Type bodyType, RequestTemplate template) throws EncodeException {
log.info("encode value:{}", object);
defaultEncoder.encode(object, bodyType, template);
}
}
@@ -0,0 +1,30 @@
package ru.demo.config;
import com.fasterxml.jackson.databind.ObjectMapper;
import feign.FeignException;
import feign.Response;
import feign.codec.Decoder;
import java.io.IOException;
import java.lang.reflect.Type;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import ru.demo.model.ClientData;
public class ResponseDecoder implements Decoder {
private static final Logger log = LoggerFactory.getLogger(ResponseDecoder.class);
private final Decoder defaultDecoder;
private final ObjectMapper mapper;
public ResponseDecoder(Decoder defaultDecoder, ObjectMapper mapper) {
this.defaultDecoder = defaultDecoder;
this.mapper = mapper;
}
@Override
public ClientData decode(Response response, Type type) throws IOException, FeignException {
var responseAsString = (String) defaultDecoder.decode(response, String.class);
log.info("response:{}", responseAsString);
return mapper.readValue(responseAsString, ClientData.class);
}
}
@@ -0,0 +1,106 @@
package ru.demo.config;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.json.JsonMapper;
import feign.Contract;
import feign.Feign;
import feign.Logger;
import feign.Retryer;
import feign.codec.Decoder;
import feign.codec.Encoder;
import feign.micrometer.MicrometerCapability;
import feign.micrometer.MicrometerObservationCapability;
import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
import io.github.resilience4j.ratelimiter.RateLimiter;
import io.github.resilience4j.ratelimiter.RateLimiterConfig;
import io.github.resilience4j.timelimiter.TimeLimiterConfig;
import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.observation.ObservationRegistry;
import java.time.Duration;
import org.springframework.boot.web.servlet.FilterRegistrationBean;
import org.springframework.cloud.circuitbreaker.resilience4j.Resilience4JCircuitBreakerFactory;
import org.springframework.cloud.circuitbreaker.resilience4j.Resilience4JConfigBuilder;
import org.springframework.cloud.client.circuitbreaker.CircuitBreaker;
import org.springframework.cloud.client.circuitbreaker.CircuitBreakerFactory;
import org.springframework.cloud.client.circuitbreaker.Customizer;
import org.springframework.cloud.openfeign.FeignClientsConfiguration;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import org.springframework.web.filter.OncePerRequestFilter;
import ru.demo.controller.ClientAdditionalInfoClient;
import ru.demo.filter.MdcFilter;
import ru.demo.metrics.MetricsManager;
import ru.demo.metrics.MicrometerMetricsManager;
@Configuration
@Import(FeignClientsConfiguration.class)
public class ServiceClientApplConf {
@Bean
public RateLimiterConfig rateLimiterConfig() {
return RateLimiterConfig.custom()
.timeoutDuration(Duration.ofMillis(100))
.limitRefreshPeriod(Duration.ofSeconds(1))
.limitForPeriod(1000)
.build();
}
@Bean
public RateLimiter rateLimiter(RateLimiterConfig config) {
return RateLimiter.of("defaultRateLimiter", config);
}
@Bean
public Customizer<Resilience4JCircuitBreakerFactory> defaultCustomizer() {
return factory -> factory.configureDefault(id -> new Resilience4JConfigBuilder(id)
.timeLimiterConfig(TimeLimiterConfig.custom()
.timeoutDuration(Duration.ofSeconds(5))
.build())
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
.build());
}
@Bean
public CircuitBreaker circuitBreaker(CircuitBreakerFactory<?, ?> circuitBreakerFactory) {
return circuitBreakerFactory.create("defaultCircuitBreaker");
}
@Bean
public FilterRegistrationBean<OncePerRequestFilter> mdcFilterRegistrationBean() {
var registrationBean = new FilterRegistrationBean<OncePerRequestFilter>();
registrationBean.setFilter(new MdcFilter());
registrationBean.setOrder(1);
return registrationBean;
}
@Bean
public ObjectMapper objectMapper() {
return JsonMapper.builder().build();
}
@Bean
public ClientAdditionalInfoClient clientAdditionalInfoClient(
Decoder decoder,
Encoder encoder,
Contract contract,
ObjectMapper mapper,
MeterRegistry meterRegistry,
ObservationRegistry observationRegistry) {
return Feign.builder()
.encoder(new RequestEncoder(encoder))
.decoder(new ResponseDecoder(decoder, mapper))
.contract(contract)
.logLevel(Logger.Level.FULL)
.addCapability(new MicrometerObservationCapability(observationRegistry)) // <-- THIS IS NEW
.addCapability(new MicrometerCapability(meterRegistry)) // <-- THIS IS NEW
.retryer(new Retryer.Default(500, 5_000, 10))
.target(ClientAdditionalInfoClient.class, "http");
}
@Bean
public MetricsManager micrometerMetricsManager(MeterRegistry meterRegistry) {
return new MicrometerMetricsManager(meterRegistry);
}
}
@@ -0,0 +1,16 @@
package ru.demo.controller;
import static ru.demo.filter.MdcFilter.HEADER_X_REQUEST_ID;
import java.net.URI;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestHeader;
import org.springframework.web.bind.annotation.RequestParam;
import ru.demo.model.ClientData;
public interface ClientAdditionalInfoClient {
@GetMapping(value = "/additional-info", consumes = "application/json")
ClientData additionalInfo(
@RequestHeader(HEADER_X_REQUEST_ID) String xRequestId, URI baseUri, @RequestParam("name") String nameVal);
}
@@ -0,0 +1,90 @@
package ru.demo.controller;
import static ru.demo.filter.MdcFilter.MDC_REQUEST_ID;
import static ru.demo.metrics.Meter.REQUEST_COUNTER;
import com.netflix.discovery.EurekaClient;
import io.github.resilience4j.core.functions.CheckedFunction;
import io.github.resilience4j.ratelimiter.RateLimiter;
import java.net.URI;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.MDC;
import org.springframework.cloud.client.circuitbreaker.CircuitBreaker;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import ru.demo.metrics.Meter;
import ru.demo.metrics.MetricsManager;
import ru.demo.model.RequestForData;
@RestController
public class ClientController {
private static final Logger log = LoggerFactory.getLogger(ClientController.class);
private final MetricsManager metricsManager;
private final ClientAdditionalInfoClient clientAdditionalInfoClient;
private final EurekaClient discoveryClient;
private final CheckedFunction<RequestForData, String> getAdditionalInfoFunction;
// curl -v -H "X-Request-Id: 123" http://localhost:8081/info?name="testClient"
public ClientController(
MetricsManager metricsManager,
ClientAdditionalInfoClient clientAdditionalInfoClient,
EurekaClient discoveryClient,
CircuitBreaker circuitBreaker,
RateLimiter rateLimiter) {
this.metricsManager = metricsManager;
this.clientAdditionalInfoClient = clientAdditionalInfoClient;
this.discoveryClient = discoveryClient;
this.getAdditionalInfoFunction = RateLimiter.decorateCheckedFunction(
rateLimiter,
requestForData -> circuitBreaker.run(() -> doRequest(requestForData), t -> {
log.error("delay call failed error:{}", t.getMessage());
return "unknown info";
}));
}
@GetMapping(value = "/info")
public String info(@RequestParam(name = "name") String name) {
var startTime = System.currentTimeMillis();
metricsManager.incrementValue(REQUEST_COUNTER);
log.info("request. name:{}", name);
String additionalInfo = null;
try {
additionalInfo = getAdditionalInfoFunction.apply(new RequestForData(name, MDC.get(MDC_REQUEST_ID)));
} catch (Throwable ex) {
log.error("can't execute additional info, name:{}, error:{}", name, ex.getMessage());
}
var requestResult = String.format("ClientInfo name:%s, additional:%s", name, additionalInfo);
var duration = System.currentTimeMillis() - startTime;
metricsManager.putValue(Meter.REQUEST_DURATION, duration);
return requestResult;
}
private String doRequest(RequestForData requestForData) {
try {
MDC.put(MDC_REQUEST_ID, requestForData.requestId());
return getAdditionalInfo(requestForData.name());
} finally {
MDC.remove(MDC_REQUEST_ID);
}
}
private String getAdditionalInfo(String name) {
try {
var clientInfo = discoveryClient.getNextServerFromEureka("SERVICE-CLIENT-INFO", false);
log.info("clientInfo from Eureka:{}", clientInfo);
var additionalInfo = clientAdditionalInfoClient.additionalInfo(
MDC.get(MDC_REQUEST_ID), new URI(clientInfo.getHomePageUrl()), name);
log.info("additionalInfo:{}", additionalInfo);
return additionalInfo.data();
} catch (Exception ex) {
log.error("can't get additional info, name:{}, error:{}", name, ex.getMessage());
return null;
}
}
}
@@ -0,0 +1,38 @@
package ru.demo.filter;
import jakarta.servlet.FilterChain;
import jakarta.servlet.ServletException;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.servlet.http.HttpServletResponse;
import java.io.IOException;
import java.util.ArrayList;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.MDC;
import org.springframework.web.filter.OncePerRequestFilter;
public class MdcFilter extends OncePerRequestFilter {
public static final String HEADER_X_REQUEST_ID = "X-Request-Id";
public static final String MDC_REQUEST_ID = "requestId";
private final Logger log = LoggerFactory.getLogger(MdcFilter.class);
@Override
protected void doFilterInternal(HttpServletRequest request, HttpServletResponse response, FilterChain filterChain)
throws ServletException, IOException {
var xRequestId = request.getHeader(HEADER_X_REQUEST_ID);
log.debug("xRequestId:{}", xRequestId);
if (xRequestId != null) {
MDC.put(MDC_REQUEST_ID, xRequestId);
}
var headerIterator = request.getHeaderNames().asIterator();
var headers = new ArrayList<String>();
while (headerIterator.hasNext()) {
headers.add(headerIterator.next());
}
log.debug("request headers:{}", headers);
response.addHeader(HEADER_X_REQUEST_ID, xRequestId);
filterChain.doFilter(request, response);
MDC.remove(MDC_REQUEST_ID);
log.debug("response headers:{}", response.getHeaderNames());
}
}
@@ -0,0 +1,16 @@
package ru.demo.metrics;
public enum Meter {
REQUEST_COUNTER("request_counter"),
REQUEST_DURATION("request_duration");
private final String meterName;
Meter(String meterName) {
this.meterName = meterName;
}
public String getMeterName() {
return meterName;
}
}
@@ -0,0 +1,11 @@
package ru.demo.metrics;
import java.util.function.Supplier;
public interface MetricsManager {
void putValue(Meter meterName, long value);
void incrementValue(Meter meterName);
void registerGauge(Meter gaugeName, Supplier<Number> gaugeGetter);
}
@@ -0,0 +1,52 @@
package ru.demo.metrics;
import io.micrometer.core.instrument.Counter;
import io.micrometer.core.instrument.Gauge;
import io.micrometer.core.instrument.MeterRegistry;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Supplier;
public class MicrometerMetricsManager implements MetricsManager {
private final MeterRegistry meterRegistry;
private final Map<String, AtomicLong> gauges = new ConcurrentHashMap<>();
private final Map<String, Counter> counters = new ConcurrentHashMap<>();
public MicrometerMetricsManager(MeterRegistry meterRegistry) {
this.meterRegistry = meterRegistry;
}
@Override
public void putValue(Meter meterName, long value) {
var gaugeName = makeMeterName(meterName);
var gauge = gauges.computeIfAbsent(gaugeName, key -> {
var newGauge = new AtomicLong();
registerGauge(meterName, newGauge::get);
return newGauge;
});
gauge.set(value);
}
@Override
public void registerGauge(Meter gaugeName, Supplier<Number> gaugeGetter) {
var builder = Gauge.builder(gaugeName.getMeterName(), gaugeGetter);
builder.register(meterRegistry);
}
@Override
public void incrementValue(Meter meterName) {
var counterName = makeMeterName(meterName);
var counter = counters.computeIfAbsent(counterName, key -> makeCounter(meterName));
counter.increment();
}
private Counter makeCounter(Meter meterName) {
var builder = Counter.builder(meterName.getMeterName());
return builder.register(meterRegistry);
}
private String makeMeterName(Meter meterName) {
return String.format("%s--", meterName.getMeterName());
}
}
@@ -0,0 +1,3 @@
package ru.demo.model;
public record ClientData(String data) {}
@@ -0,0 +1,3 @@
package ru.demo.model;
public record RequestForData(String name, String requestId) {}
@@ -0,0 +1,18 @@
spring:
application:
name: service-client
cloud:
config:
fail-fast: true
retry:
initial-interval: 5000
max-attempts: 10
max-interval: 5000
multiplier: 1.2
config:
import: optional:configserver:http://localhost:8888
codec:
max-in-memory-size: 10MB
@@ -0,0 +1,26 @@
<configuration scan="true" scanPeriod="10 seconds">
<jmxConfigurator />
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - [%X{requestId}] %msg%n</pattern>
</encoder>
</appender>
<appender name="kafka" class="ru.appender.kafka.LogAppender">
<bootstrapServers>localhost:9092</bootstrapServers>
<topicName>applLogs</topicName>
<encoder class="net.logstash.logback.encoder.LogstashEncoder">
<fieldNames>
<version>[ignore]</version>
</fieldNames>
<customFields>{"appname":"service-client"}</customFields>
<includeMdcKeyName>requestId</includeMdcKeyName>
</encoder>
</appender>
<root level="info">
<appender-ref ref="STDOUT" />
<appender-ref ref="kafka" />
</root>
</configuration>
@@ -0,0 +1,39 @@
plugins {
id("com.google.cloud.tools.jib")
}
dependencies {
implementation(project(":kafka-log-appender"))
implementation("net.logstash.logback:logstash-logback-encoder")
implementation ("org.springframework.boot:spring-boot-starter-web")
implementation("org.springframework.boot:spring-boot-starter-actuator")
implementation("io.micrometer:micrometer-registry-prometheus")
implementation("org.springframework.cloud:spring-cloud-starter-config")
implementation("org.springframework.cloud:spring-cloud-starter-netflix-eureka-client")
implementation("io.micrometer:micrometer-tracing-bridge-otel") // bridges the Micrometer Observation API to OpenTelemetry.
implementation("io.opentelemetry:opentelemetry-exporter-zipkin") // reports traces to Zipkin.
}
jib {
container {
creationTime.set("USE_CURRENT_TIMESTAMP")
}
from {
image = "bellsoft/liberica-openjdk-alpine-musl:21.0.1"
}
to {
image = "localrun/service-order"
tags = setOf(project.version.toString())
}
}
tasks {
build {
dependsOn(spotlessApply)
dependsOn(jibBuildTar)
}
}
@@ -0,0 +1,27 @@
#!/bin/bash
../gradlew :service-order:build
docker load --input build/jib-image.tar
docker stop service-order-1
docker stop service-order-2
docker run --rm -d --name service-order-1 \
--memory=256m \
--cpus 1 \
--network="host" \
-e JAVA_TOOL_OPTIONS="-XX:InitialRAMPercentage=80 -XX:MaxRAMPercentage=80" \
-e SPRING_APPLICATION_INSTANCE_ID="i1" \
-e SERVER_PORT=8091 \
localrun/service-order:latest
docker run --rm -d --name service-order-2 \
--memory=256m \
--cpus 1 \
--network="host" \
-e JAVA_TOOL_OPTIONS="-XX:InitialRAMPercentage=80 -XX:MaxRAMPercentage=80" \
-e SPRING_APPLICATION_INSTANCE_ID="i2" \
-e SERVER_PORT=8092 \
localrun/service-order:latest
@@ -0,0 +1,15 @@
package ru.demo;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import ru.demo.controller.InstanceId;
@Configuration
public class Config {
@Bean
InstanceId instanceId(@Value("${spring.application.instance_id}") String id) {
return new InstanceId(id);
}
}
@@ -0,0 +1,11 @@
package ru.demo;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.builder.SpringApplicationBuilder;
@SpringBootApplication
public class ServiceOrder {
public static void main(String[] args) {
new SpringApplicationBuilder().sources(ServiceOrder.class).run(args);
}
}
@@ -0,0 +1,3 @@
package ru.demo.controller;
public record InstanceId(String name) {}
@@ -0,0 +1,28 @@
package ru.demo.controller;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class OrderController {
private static final Logger logger = LoggerFactory.getLogger(OrderController.class);
private final InstanceId instanceId;
public OrderController(InstanceId instanceId) {
this.instanceId = instanceId;
}
// curl -v http://localhost:8082/info?id="idClient"
// curl -v http://localhost:8091/info?id="idClient"
// curl -v http://localhost:8092/info?id="idClient"
@GetMapping(value = "/info")
public String info(@RequestParam(name = "id") String id) {
logger.info("instanceId:{}, request. id:{}", instanceId.name(), id);
return String.format("%s, Order id:%s", instanceId.name(), id);
}
}
@@ -0,0 +1,18 @@
spring:
application:
name: service-order
instance_id: i0
cloud:
config:
fail-fast: true
retry:
initial-interval: 5000
max-attempts: 10
max-interval: 5000
multiplier: 1.2
config:
import: optional:configserver:http://localhost:8888
codec:
max-in-memory-size: 10MB
@@ -0,0 +1,26 @@
<configuration scan="true" scanPeriod="10 seconds">
<jmxConfigurator />
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - [%X{requestId}] %msg%n</pattern>
</encoder>
</appender>
<appender name="kafka" class="ru.appender.kafka.LogAppender">
<bootstrapServers>localhost:9092</bootstrapServers>
<topicName>applLogs</topicName>
<encoder class="net.logstash.logback.encoder.LogstashEncoder">
<fieldNames>
<version>[ignore]</version>
</fieldNames>
<customFields>{"appname":"service-order"}</customFields>
<includeMdcKeyName>requestId</includeMdcKeyName>
</encoder>
</appender>
<root level="info">
<appender-ref ref="STDOUT" />
<appender-ref ref="kafka" />
</root>
</configuration>
+27
View File
@@ -0,0 +1,27 @@
pluginManagement {
val jgitver: String by settings
val dependencyManagement: String by settings
val springframeworkBoot: String by settings
val sonarlint: String by settings
val spotless: String by settings
val jib: String by settings
plugins {
id("fr.brouillard.oss.gradle.jgitver") version jgitver
id("io.spring.dependency-management") version dependencyManagement
id("org.springframework.boot") version springframeworkBoot
id("name.remal.sonarlint") version sonarlint
id("com.diffplug.spotless") version spotless
id("com.google.cloud.tools.jib") version jib
}
}
rootProject.name = "spring-cloud"
include("api-gateway")
include("config-server")
include("eureka-server")
include("service-client")
include("service-client-info")
include("service-order")
include("kafka-log-appender")