mirror of
https://github.com/iflytek/skillhub.git
synced 2026-10-06 02:48:28 +00:00
Merge issue #597 observability fix into big-main
Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com>
This commit is contained in:
commit
98e2bd7bf6
6 changed files with 130 additions and 5 deletions
|
|
@ -1,14 +1,16 @@
|
|||
package com.iflytek.skillhub.config;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executor;
|
||||
import java.util.concurrent.ThreadPoolExecutor;
|
||||
import org.slf4j.MDC;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.core.task.TaskDecorator;
|
||||
import org.springframework.scheduling.annotation.EnableAsync;
|
||||
import org.springframework.scheduling.annotation.EnableScheduling;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
|
||||
import java.util.concurrent.Executor;
|
||||
import java.util.concurrent.ThreadPoolExecutor;
|
||||
|
||||
/**
|
||||
* Enables asynchronous event handling and other background execution features used by the
|
||||
* application module.
|
||||
|
|
@ -26,7 +28,30 @@ public class AsyncConfig {
|
|||
executor.setQueueCapacity(100);
|
||||
executor.setThreadNamePrefix("skillhub-event-");
|
||||
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
|
||||
executor.setTaskDecorator(mdcTaskDecorator());
|
||||
executor.initialize();
|
||||
return executor;
|
||||
}
|
||||
|
||||
private TaskDecorator mdcTaskDecorator() {
|
||||
return task -> {
|
||||
Map<String, String> callerContext = MDC.getCopyOfContextMap();
|
||||
return () -> {
|
||||
Map<String, String> executorContext = MDC.getCopyOfContextMap();
|
||||
try {
|
||||
restoreMdc(callerContext);
|
||||
task.run();
|
||||
} finally {
|
||||
restoreMdc(executorContext);
|
||||
}
|
||||
};
|
||||
};
|
||||
}
|
||||
|
||||
private void restoreMdc(Map<String, String> context) {
|
||||
MDC.clear();
|
||||
if (context != null) {
|
||||
MDC.setContextMap(context);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -60,4 +60,8 @@ public class SkillHubMetrics {
|
|||
"operation", operation
|
||||
).increment();
|
||||
}
|
||||
|
||||
public void incrementSearchRebuildFailure() {
|
||||
meterRegistry.counter("skillhub.search.rebuild.failure").increment();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
package com.iflytek.skillhub.service;
|
||||
|
||||
import com.iflytek.skillhub.metrics.SkillHubMetrics;
|
||||
import com.iflytek.skillhub.search.SearchRebuildService;
|
||||
import java.util.List;
|
||||
import org.slf4j.Logger;
|
||||
|
|
@ -15,9 +16,12 @@ public class LabelSearchSyncService {
|
|||
private static final Logger log = LoggerFactory.getLogger(LabelSearchSyncService.class);
|
||||
|
||||
private final SearchRebuildService searchRebuildService;
|
||||
private final SkillHubMetrics metrics;
|
||||
|
||||
public LabelSearchSyncService(SearchRebuildService searchRebuildService) {
|
||||
public LabelSearchSyncService(SearchRebuildService searchRebuildService,
|
||||
SkillHubMetrics metrics) {
|
||||
this.searchRebuildService = searchRebuildService;
|
||||
this.metrics = metrics;
|
||||
}
|
||||
|
||||
@Async("skillhubEventExecutor")
|
||||
|
|
@ -25,6 +29,7 @@ public class LabelSearchSyncService {
|
|||
try {
|
||||
searchRebuildService.rebuildBySkill(skillId);
|
||||
} catch (RuntimeException ex) {
|
||||
metrics.incrementSearchRebuildFailure();
|
||||
log.error("Failed to rebuild search document for skill {}", skillId, ex);
|
||||
}
|
||||
}
|
||||
|
|
@ -41,6 +46,7 @@ public class LabelSearchSyncService {
|
|||
try {
|
||||
searchRebuildService.rebuildBySkill(skillId);
|
||||
} catch (RuntimeException ex) {
|
||||
metrics.incrementSearchRebuildFailure();
|
||||
log.error("Failed to rebuild search document for skill {} after label change", skillId, ex);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2,9 +2,13 @@ package com.iflytek.skillhub.config;
|
|||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.MDC;
|
||||
import org.springframework.scheduling.annotation.EnableAsync;
|
||||
import org.springframework.scheduling.annotation.EnableScheduling;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
|
||||
class AsyncConfigTest {
|
||||
|
||||
|
|
@ -13,4 +17,26 @@ class AsyncConfigTest {
|
|||
assertThat(AsyncConfig.class).hasAnnotation(EnableAsync.class);
|
||||
assertThat(AsyncConfig.class).hasAnnotation(EnableScheduling.class);
|
||||
}
|
||||
|
||||
@Test
|
||||
void skillhubEventExecutor_propagatesAndClearsMdc() throws Exception {
|
||||
ThreadPoolTaskExecutor executor =
|
||||
(ThreadPoolTaskExecutor) new AsyncConfig().skillhubEventExecutor();
|
||||
try {
|
||||
MDC.put("requestId", "req-597");
|
||||
CompletableFuture<String> propagatedRequestId = new CompletableFuture<>();
|
||||
executor.execute(() -> propagatedRequestId.complete(MDC.get("requestId")));
|
||||
MDC.clear();
|
||||
|
||||
assertThat(propagatedRequestId.get(5, TimeUnit.SECONDS)).isEqualTo("req-597");
|
||||
|
||||
CompletableFuture<String> nextRequestId = new CompletableFuture<>();
|
||||
executor.execute(() -> nextRequestId.complete(MDC.get("requestId")));
|
||||
|
||||
assertThat(nextRequestId.get(5, TimeUnit.SECONDS)).isNull();
|
||||
} finally {
|
||||
MDC.clear();
|
||||
executor.shutdown();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -39,6 +39,7 @@ class PrometheusEndpointTest {
|
|||
skillHubMetrics.incrementUserRegister();
|
||||
skillHubMetrics.recordLocalLogin(true);
|
||||
skillHubMetrics.incrementSkillPublish("global", "PENDING_REVIEW");
|
||||
skillHubMetrics.incrementSearchRebuildFailure();
|
||||
|
||||
assertThat(environment.getProperty("management.endpoints.web.exposure.include"))
|
||||
.doesNotContain("prometheus")
|
||||
|
|
@ -54,5 +55,8 @@ class PrometheusEndpointTest {
|
|||
.tag("status", "PENDING_REVIEW")
|
||||
.counter()
|
||||
.count()).isEqualTo(1.0d);
|
||||
assertThat(meterRegistry.get("skillhub.search.rebuild.failure")
|
||||
.counter()
|
||||
.count()).isEqualTo(1.0d);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,13 +1,19 @@
|
|||
package com.iflytek.skillhub.service;
|
||||
|
||||
import com.iflytek.skillhub.metrics.SkillHubMetrics;
|
||||
import com.iflytek.skillhub.search.SearchRebuildService;
|
||||
import io.micrometer.core.instrument.simple.SimpleMeterRegistry;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.Mockito.doThrow;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.verifyNoInteractions;
|
||||
import static org.mockito.Mockito.verifyNoMoreInteractions;
|
||||
|
||||
class LabelSearchSyncServiceTest {
|
||||
|
|
@ -15,7 +21,8 @@ class LabelSearchSyncServiceTest {
|
|||
@Test
|
||||
void rebuildSkillsShouldSkipNullsAndDuplicatesWhileProcessingLargeLists() {
|
||||
SearchRebuildService rebuildService = mock(SearchRebuildService.class);
|
||||
LabelSearchSyncService service = new LabelSearchSyncService(rebuildService);
|
||||
SkillHubMetrics metrics = mock(SkillHubMetrics.class);
|
||||
LabelSearchSyncService service = new LabelSearchSyncService(rebuildService, metrics);
|
||||
List<Long> skillIds = new ArrayList<>();
|
||||
skillIds.add(null);
|
||||
for (long i = 1; i <= 120; i++) {
|
||||
|
|
@ -30,5 +37,58 @@ class LabelSearchSyncServiceTest {
|
|||
verify(rebuildService).rebuildBySkill(i);
|
||||
}
|
||||
verifyNoMoreInteractions(rebuildService);
|
||||
verifyNoInteractions(metrics);
|
||||
}
|
||||
|
||||
@Test
|
||||
void rebuildSkillFailureShouldIncrementMetric() {
|
||||
SearchRebuildService rebuildService = mock(SearchRebuildService.class);
|
||||
doThrow(new IllegalStateException("search unavailable"))
|
||||
.when(rebuildService)
|
||||
.rebuildBySkill(42L);
|
||||
|
||||
contextRunner(rebuildService).run(context -> {
|
||||
LabelSearchSyncService service = context.getBean(LabelSearchSyncService.class);
|
||||
SimpleMeterRegistry meterRegistry = context.getBean(SimpleMeterRegistry.class);
|
||||
|
||||
service.rebuildSkill(42L);
|
||||
|
||||
assertThat(meterRegistry.get("skillhub.search.rebuild.failure").counter().count())
|
||||
.isEqualTo(1.0d);
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void rebuildSkillsShouldCountEachFailureAndContinue() {
|
||||
SearchRebuildService rebuildService = mock(SearchRebuildService.class);
|
||||
doThrow(new IllegalStateException("search unavailable"))
|
||||
.when(rebuildService)
|
||||
.rebuildBySkill(2L);
|
||||
doThrow(new IllegalStateException("search unavailable"))
|
||||
.when(rebuildService)
|
||||
.rebuildBySkill(3L);
|
||||
|
||||
contextRunner(rebuildService).run(context -> {
|
||||
LabelSearchSyncService service = context.getBean(LabelSearchSyncService.class);
|
||||
SimpleMeterRegistry meterRegistry = context.getBean(SimpleMeterRegistry.class);
|
||||
|
||||
service.rebuildSkills(List.of(1L, 2L, 3L, 4L));
|
||||
|
||||
assertThat(meterRegistry.get("skillhub.search.rebuild.failure").counter().count())
|
||||
.isEqualTo(2.0d);
|
||||
verify(rebuildService).rebuildBySkill(1L);
|
||||
verify(rebuildService).rebuildBySkill(2L);
|
||||
verify(rebuildService).rebuildBySkill(3L);
|
||||
verify(rebuildService).rebuildBySkill(4L);
|
||||
verifyNoMoreInteractions(rebuildService);
|
||||
});
|
||||
}
|
||||
|
||||
private ApplicationContextRunner contextRunner(SearchRebuildService rebuildService) {
|
||||
return new ApplicationContextRunner()
|
||||
.withBean(SearchRebuildService.class, () -> rebuildService)
|
||||
.withBean(SimpleMeterRegistry.class)
|
||||
.withBean(SkillHubMetrics.class)
|
||||
.withBean(LabelSearchSyncService.class);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue