diff --git a/.github/workflows/backend-cd.yml b/.github/workflows/backend-cd.yml index 57c70e2a..607d4751 100644 --- a/.github/workflows/backend-cd.yml +++ b/.github/workflows/backend-cd.yml @@ -60,6 +60,7 @@ jobs: if: steps.changes.outputs.deploy == 'true' env: DOCKERHUB_USERNAME: ${{ secrets.DOCKERHUB_USERNAME }} + REDIS_PASSWORD: ${{ secrets.REDIS_PASSWORD }} run: | cd infra/docker docker compose -f docker-compose-prod.yml build @@ -70,18 +71,20 @@ jobs: uses: appleboy/ssh-action@v1.0.3 env: DOCKERHUB_USERNAME: ${{ secrets.DOCKERHUB_USERNAME }} + REDIS_PASSWORD: ${{ secrets.REDIS_PASSWORD }} DEPLOY_SHA: ${{ github.event.workflow_run.head_sha }} with: host: ${{ secrets.EC2_HOST }} username: ${{ secrets.EC2_USER }} key: ${{ secrets.EC2_SSH_KEY }} - envs: DOCKERHUB_USERNAME,DEPLOY_SHA + envs: DOCKERHUB_USERNAME,REDIS_PASSWORD,DEPLOY_SHA script: | cd ~/CoinFlow git fetch origin master git checkout "$DEPLOY_SHA" export DOCKERHUB_USERNAME=$DOCKERHUB_USERNAME + export REDIS_PASSWORD=$REDIS_PASSWORD cd infra/docker docker compose -f docker-compose-prod.yml pull diff --git a/backend/coinflow-collector-app/Dockerfile b/backend/coinflow-collector-app/Dockerfile index fcbb2bf9..8a815b2d 100644 --- a/backend/coinflow-collector-app/Dockerfile +++ b/backend/coinflow-collector-app/Dockerfile @@ -13,7 +13,7 @@ COPY build/libs/*-SNAPSHOT.jar app.jar ENV PROFILE=prod # 포트 개방 (Collector 모듈 내부 포트) -EXPOSE 8083 +EXPOSE 8082 # 메모리 최적화 옵션 및 실행 ENTRYPOINT ["sh", "-c", "java -Dspring.profiles.active=${PROFILE} -Xms256m -Xmx512m -jar app.jar"] diff --git a/backend/coinflow-collector-app/src/main/resources/application-collector.yml b/backend/coinflow-collector-app/src/main/resources/application-collector.yml index 7e2c3980..a80d9b67 100644 --- a/backend/coinflow-collector-app/src/main/resources/application-collector.yml +++ b/backend/coinflow-collector-app/src/main/resources/application-collector.yml @@ -22,3 +22,9 @@ management: metrics: tags: application: coinflow-collector + +redis: + stream: + tick: + stream-key: ${REDIS_STREAM_TICK_STREAMKEY:tick:raw} + max-length: ${REDIS_STREAM_TICK_MAXLENGTH:200000} diff --git a/backend/coinflow-common/src/main/java/com/coinflow/monitoring/constant/MetricConstants.java b/backend/coinflow-common/src/main/java/com/coinflow/monitoring/constant/MetricConstants.java index 07251c86..1efbbb81 100644 --- a/backend/coinflow-common/src/main/java/com/coinflow/monitoring/constant/MetricConstants.java +++ b/backend/coinflow-common/src/main/java/com/coinflow/monitoring/constant/MetricConstants.java @@ -10,12 +10,15 @@ private MetricConstants() { public static final String STREAM_ACK_COUNT = "stream.ack.count"; public static final String STREAM_ACK_LATENCY = "stream.ack.latency"; public static final String STREAM_BACKLOG_COUNT = "stream.backlog.count"; + public static final String STREAM_BACKLOG_RETENTION_RATIO = "stream.backlog.retention.ratio"; + public static final String STREAM_RETENTION_WARNING_COUNT = "stream.retention.warning.count"; public static final String STREAM_PEL_COUNT = "stream.pel.count"; public static final String REDIS_COMMAND_COUNT = "redis.command.count"; // Collector: 유입량 및 발행 지표 public static final String WEBSOCKET_RECEIVE_COUNT = "tick.receive.count"; public static final String STREAM_PUBLISH_LATENCY = "stream.publish.latency"; + public static final String STREAM_PUBLISH_FAILURE_COUNT = "stream.publish.failure.count"; // Consumer: 틱 처리 전체 지표 public static final String TICK_PROCESS_LATENCY = "tick.process.latency"; @@ -36,10 +39,9 @@ private MetricConstants() { public static final String VALUE_SUCCESS = "success"; public static final String VALUE_FAILURE = "failure"; public static final String VALUE_MODULE_CONSUMER = "consumer"; + public static final String VALUE_MODULE_COLLECTOR = "collector"; public static final String VALUE_NA = "NA"; public static final String VALUE_FLUSH_SIZE = "size"; public static final String VALUE_FLUSH_INTERVAL = "interval"; - // Redis Stream Configuration - public static final long STREAM_MAX_LEN = 1_000_000L; } diff --git a/backend/coinflow-consumer-app/Dockerfile b/backend/coinflow-consumer-app/Dockerfile index 4229c116..f670cba2 100644 --- a/backend/coinflow-consumer-app/Dockerfile +++ b/backend/coinflow-consumer-app/Dockerfile @@ -13,7 +13,7 @@ COPY build/libs/*-SNAPSHOT.jar app.jar ENV PROFILE=prod # 포트 개방 (Consumer 모듈 내부 포트) -EXPOSE 8082 +EXPOSE 8081 # 메모리 최적화 옵션 및 실행 ENTRYPOINT ["sh", "-c", "java -Dspring.profiles.active=${PROFILE} -Xms512m -Xmx1024m -jar app.jar"] diff --git a/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/ConsumerApplicationShutdown.java b/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/ConsumerApplicationShutdown.java new file mode 100644 index 00000000..e152e5d3 --- /dev/null +++ b/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/ConsumerApplicationShutdown.java @@ -0,0 +1,25 @@ +package com.coinflow.config; + +import java.util.concurrent.atomic.AtomicBoolean; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.stereotype.Component; + +@Component +@RequiredArgsConstructor +@Slf4j +public class ConsumerApplicationShutdown { + + private final ConfigurableApplicationContext applicationContext; + private final AtomicBoolean shutdownRequested = new AtomicBoolean(false); + + public void request() { + if (!shutdownRequested.compareAndSet(false, true)) { + return; + } + + log.error("Closing consumer application after a fatal Redis Stream subscription error"); + applicationContext.close(); + } +} diff --git a/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/RedisConsumerConfig.java b/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/RedisConsumerConfig.java index f8627612..91dd13de 100644 --- a/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/RedisConsumerConfig.java +++ b/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/RedisConsumerConfig.java @@ -4,26 +4,22 @@ import com.coinflow.consumer.TickRawEventConsumer; import jakarta.annotation.PreDestroy; import java.time.Duration; -import java.util.Map; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.redis.connection.RedisConnectionFactory; -import org.springframework.data.redis.connection.RedisStreamCommands.XAddOptions; import org.springframework.data.redis.connection.stream.Consumer; import org.springframework.data.redis.connection.stream.MapRecord; import org.springframework.data.redis.connection.stream.ReadOffset; import org.springframework.data.redis.connection.stream.StreamOffset; -import org.springframework.data.redis.connection.stream.StreamRecords; -import org.springframework.data.redis.core.RedisTemplate; import org.springframework.data.redis.serializer.RedisSerializer; import org.springframework.data.redis.stream.StreamMessageListenerContainer; +import org.springframework.data.redis.stream.StreamMessageListenerContainer.StreamReadRequest; import org.springframework.data.redis.stream.StreamMessageListenerContainer.StreamMessageListenerContainerOptions; -import static com.coinflow.monitoring.constant.MetricConstants.STREAM_MAX_LEN; - /** * Redis Stream 소비자(Consumer) 설정을 담당하며, 바이너리 수신(Phase 3.1)을 지원합니다. */ @@ -33,23 +29,17 @@ @RequiredArgsConstructor public class RedisConsumerConfig { - private static final String ERROR_BUSYGROUP = "BUSYGROUP"; - private static final String ERROR_NO_SUCH_KEY = "No such key"; - private static final String ERROR_NO_GROUP = "NOGROUP"; - private static final String DUMMY_EVENT_KEY = "init-event"; - private static final String DUMMY_EVENT_VALUE = "true"; - private final RedisConnectionFactory connectionFactory; private final TickRawEventConsumer consumer; private final TickConsumerProperties properties; + private final RedisConsumerGroupManager consumerGroupManager; private StreamMessageListenerContainer> container; @Bean - public StreamMessageListenerContainer> tickStreamContainer( - RedisTemplate redisTemplate) { - - initializeConsumerGroup(redisTemplate); + @ConditionalOnProperty(prefix = "redis.stream.tick", name = "enabled", havingValue = "true", matchIfMissing = true) + public StreamMessageListenerContainer> tickStreamContainer() { + consumerGroupManager.ensureConsumerGroup(); // 바이너리 수신을 위한 컨테이너 옵션 설정 (ByteArrayRedisSerializer) @SuppressWarnings("unchecked") @@ -63,43 +53,19 @@ public StreamMessageListenerContainer> .build(); container = StreamMessageListenerContainer.create(connectionFactory, options); - container.receive( - Consumer.from(properties.group(), properties.consumerName()), - StreamOffset.create(properties.streamKey(), ReadOffset.lastConsumed()), - this.consumer); + StreamReadRequest readRequest = StreamReadRequest + .builder(StreamOffset.create(properties.streamKey(), ReadOffset.lastConsumed())) + .consumer(Consumer.from(properties.group(), properties.consumerName())) + .errorHandler(consumerGroupManager::handleSubscriptionError) + .cancelOnError(consumerGroupManager::shouldCancelSubscription) + .build(); + container.register(readRequest, consumer); container.start(); log.info("Successfully started Redis Stream Container (Binary Mode) for group: {}", properties.group()); return container; } - private void initializeConsumerGroup(RedisTemplate redisTemplate) { - String streamKey = properties.streamKey(); - String group = properties.group(); - - try { - redisTemplate.opsForStream().createGroup(streamKey, ReadOffset.latest(), group); - } catch (Exception e) { - String msg = e.getMessage() != null ? e.getMessage() : ""; - if (msg.contains(ERROR_BUSYGROUP)) { - log.info("Redis consumer group already exists: {}", group); - } else if (msg.contains(ERROR_NO_SUCH_KEY) || msg.contains(ERROR_NO_GROUP)) { - log.warn("Redis stream does not exist. Initializing stream and group: {}", group); - - // 더미 메시지 발행 시에도 MAXLEN 적용 (안정성 강화) - XAddOptions options = XAddOptions.maxlen(STREAM_MAX_LEN).approximateTrimming(true); - MapRecord record = StreamRecords.newRecord() - .in(streamKey) - .ofMap(Map.of(DUMMY_EVENT_KEY, DUMMY_EVENT_VALUE)); - - redisTemplate.opsForStream().add(record, options); - redisTemplate.opsForStream().createGroup(streamKey, ReadOffset.latest(), group); - } else { - log.error("Critical error during Redis Consumer Group initialization: {}", msg); - } - } - } - @PreDestroy public void shutdown() { if (container != null && container.isRunning()) { diff --git a/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/RedisConsumerGroupManager.java b/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/RedisConsumerGroupManager.java new file mode 100644 index 00000000..fe027dcf --- /dev/null +++ b/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/RedisConsumerGroupManager.java @@ -0,0 +1,89 @@ +package com.coinflow.config; + +import com.coinflow.config.properties.TickConsumerProperties; +import java.nio.charset.StandardCharsets; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.data.redis.connection.stream.ReadOffset; +import org.springframework.data.redis.core.RedisCallback; +import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.stereotype.Component; + +@Component +@RequiredArgsConstructor +@Slf4j +public class RedisConsumerGroupManager { + + private static final String ERROR_BUSY_GROUP = "BUSYGROUP"; + private static final String ERROR_NO_GROUP = "NOGROUP"; + + private final RedisTemplate redisTemplate; + private final TickConsumerProperties properties; + private final ConsumerApplicationShutdown applicationShutdown; + + public void ensureConsumerGroup() { + try { + redisTemplate.execute((RedisCallback) connection -> + connection.streamCommands().xGroupCreate( + raw(properties.streamKey()), + properties.group(), + ReadOffset.latest(), + true)); + log.info("Created Redis consumer group. stream={}, group={}", + properties.streamKey(), properties.group()); + } catch (RuntimeException e) { + if (containsError(e, ERROR_BUSY_GROUP)) { + log.info("Redis consumer group already exists. stream={}, group={}", + properties.streamKey(), properties.group()); + return; + } + throw new IllegalStateException( + "Failed to initialize Redis consumer group. stream=" + properties.streamKey() + + ", group=" + properties.group(), + e); + } + } + + public void handleSubscriptionError(Throwable error) { + if (!isNoGroup(error)) { + log.error("Redis Stream subscription failed. stream={}, group={}", + properties.streamKey(), properties.group(), error); + applicationShutdown.request(); + return; + } + + log.warn("Redis consumer group is missing. Recreating group without cancelling subscription. " + + "stream={}, group={}", + properties.streamKey(), properties.group()); + try { + ensureConsumerGroup(); + } catch (RuntimeException recoveryError) { + log.error("Failed to recover missing Redis consumer group. stream={}, group={}", + properties.streamKey(), properties.group(), recoveryError); + } + } + + public boolean shouldCancelSubscription(Throwable error) { + return !isNoGroup(error); + } + + boolean isNoGroup(Throwable error) { + return containsError(error, ERROR_NO_GROUP); + } + + private static boolean containsError(Throwable error, String code) { + Throwable current = error; + while (current != null) { + String message = current.getMessage(); + if (message != null && message.contains(code)) { + return true; + } + current = current.getCause(); + } + return false; + } + + private static byte[] raw(String value) { + return value.getBytes(StandardCharsets.UTF_8); + } +} diff --git a/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/properties/TickConsumerProperties.java b/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/properties/TickConsumerProperties.java index d728c250..74e4a3aa 100644 --- a/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/properties/TickConsumerProperties.java +++ b/backend/coinflow-consumer-app/src/main/java/com/coinflow/config/properties/TickConsumerProperties.java @@ -1,6 +1,9 @@ package com.coinflow.config.properties; +import jakarta.validation.constraints.DecimalMax; +import jakarta.validation.constraints.DecimalMin; import jakarta.validation.constraints.NotBlank; +import jakarta.validation.constraints.Positive; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.validation.annotation.Validated; @@ -9,6 +12,8 @@ public record TickConsumerProperties( @NotBlank String streamKey, @NotBlank String group, - @NotBlank String consumerName + @NotBlank String consumerName, + @Positive long maxLength, + @DecimalMin("0.0") @DecimalMax("1.0") double lagWarningRatio ) { } diff --git a/backend/coinflow-consumer-app/src/main/java/com/coinflow/monitoring/PelRecoveryWorker.java b/backend/coinflow-consumer-app/src/main/java/com/coinflow/monitoring/PelRecoveryWorker.java index 9bc02ca0..6f82e85a 100644 --- a/backend/coinflow-consumer-app/src/main/java/com/coinflow/monitoring/PelRecoveryWorker.java +++ b/backend/coinflow-consumer-app/src/main/java/com/coinflow/monitoring/PelRecoveryWorker.java @@ -96,8 +96,8 @@ public void monitorPel() { pelCountGauge.set((double) overThresholdCount); } catch (Exception e) { - log.error("Failed to monitor Redis stream PEL. stream={}, group={}, error={}", - streamKey, group, e.getMessage()); + log.error("Failed to monitor Redis stream PEL. stream={}, group={}", + streamKey, group, e); } } } diff --git a/backend/coinflow-consumer-app/src/main/java/com/coinflow/monitoring/StreamLagMonitorWorker.java b/backend/coinflow-consumer-app/src/main/java/com/coinflow/monitoring/StreamLagMonitorWorker.java index 008a24cc..bead0c89 100644 --- a/backend/coinflow-consumer-app/src/main/java/com/coinflow/monitoring/StreamLagMonitorWorker.java +++ b/backend/coinflow-consumer-app/src/main/java/com/coinflow/monitoring/StreamLagMonitorWorker.java @@ -3,6 +3,7 @@ import com.coinflow.config.properties.TickConsumerProperties; import io.micrometer.core.instrument.Counter; import jakarta.annotation.PostConstruct; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -13,6 +14,8 @@ import static com.coinflow.monitoring.constant.MetricConstants.REDIS_COMMAND_COUNT; import static com.coinflow.monitoring.constant.MetricConstants.STREAM_BACKLOG_COUNT; +import static com.coinflow.monitoring.constant.MetricConstants.STREAM_BACKLOG_RETENTION_RATIO; +import static com.coinflow.monitoring.constant.MetricConstants.STREAM_RETENTION_WARNING_COUNT; import static com.coinflow.monitoring.constant.MetricConstants.TAG_COMMAND; import static com.coinflow.monitoring.constant.MetricConstants.TAG_FLUSH_REASON; import static com.coinflow.monitoring.constant.MetricConstants.TAG_MODULE; @@ -33,11 +36,15 @@ public class StreamLagMonitorWorker { private final MetricRecorder metricRecorder; private AtomicReference backlogGauge; + private AtomicReference backlogRetentionRatioGauge; private Counter xinfoCounter; + private final AtomicBoolean retentionWarningActive = new AtomicBoolean(false); @PostConstruct public void init() { this.backlogGauge = metricRecorder.registerGauge(STREAM_BACKLOG_COUNT, 0.0, TAG_MODULE, VALUE_MODULE_CONSUMER); + this.backlogRetentionRatioGauge = metricRecorder.registerGauge( + STREAM_BACKLOG_RETENTION_RATIO, 0.0, TAG_MODULE, VALUE_MODULE_CONSUMER); this.xinfoCounter = metricRecorder.getCounter(REDIS_COMMAND_COUNT, TAG_COMMAND, "XINFO", TAG_FLUSH_REASON, VALUE_NA); @@ -64,9 +71,7 @@ public void monitorLag() { Object lagObj = g.getRaw().get("lag"); if (lagObj instanceof Number) { double lagValue = ((Number) lagObj).doubleValue(); - backlogGauge.set(lagValue); - log.debug("Redis Stream Lag monitored: stream={}, group={}, lag={}", - streamKey, group, lagValue); + updateBacklogMetrics(lagValue); } else { log.warn("Stream lag information is not available for group: {}. Please check Redis version (7.0+ required).", group); } @@ -77,4 +82,32 @@ public void monitorLag() { streamKey, group, e.getMessage()); } } + + void updateBacklogMetrics(double lagValue) { + double retentionRatio = lagValue / properties.maxLength(); + backlogGauge.set(lagValue); + backlogRetentionRatioGauge.set(retentionRatio); + + if (retentionRatio >= properties.lagWarningRatio()) { + if (retentionWarningActive.compareAndSet(false, true)) { + metricRecorder.increment( + STREAM_RETENTION_WARNING_COUNT, + TAG_MODULE, + VALUE_MODULE_CONSUMER); + log.error("Redis Stream retention warning. stream={}, group={}, lag={}, maxLength={}, ratio={}", + properties.streamKey(), properties.group(), lagValue, + properties.maxLength(), retentionRatio); + } + return; + } + + if (retentionWarningActive.compareAndSet(true, false)) { + log.info("Redis Stream retention recovered. stream={}, group={}, lag={}, maxLength={}, ratio={}", + properties.streamKey(), properties.group(), lagValue, + properties.maxLength(), retentionRatio); + } else { + log.debug("Redis Stream Lag monitored: stream={}, group={}, lag={}, ratio={}", + properties.streamKey(), properties.group(), lagValue, retentionRatio); + } + } } diff --git a/backend/coinflow-consumer-app/src/main/resources/application-consumer.yml b/backend/coinflow-consumer-app/src/main/resources/application-consumer.yml index 9090e889..f245b36f 100644 --- a/backend/coinflow-consumer-app/src/main/resources/application-consumer.yml +++ b/backend/coinflow-consumer-app/src/main/resources/application-consumer.yml @@ -30,6 +30,8 @@ redis: stream-key: tick:raw group: tick-consumer-group consumer-name: consumer-1 + max-length: ${REDIS_STREAM_TICK_MAXLENGTH:200000} + lag-warning-ratio: ${REDIS_STREAM_TICK_LAGWARNINGRATIO:0.8} coinflow: async: diff --git a/backend/coinflow-consumer-app/src/test/java/com/coinflow/config/ConsumerApplicationShutdownTest.java b/backend/coinflow-consumer-app/src/test/java/com/coinflow/config/ConsumerApplicationShutdownTest.java new file mode 100644 index 00000000..f9e222cc --- /dev/null +++ b/backend/coinflow-consumer-app/src/test/java/com/coinflow/config/ConsumerApplicationShutdownTest.java @@ -0,0 +1,22 @@ +package com.coinflow.config; + +import org.junit.jupiter.api.Test; +import org.springframework.context.ConfigurableApplicationContext; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +class ConsumerApplicationShutdownTest { + + @Test + void closesApplicationContextOnlyOnce() { + ConfigurableApplicationContext applicationContext = mock(ConfigurableApplicationContext.class); + ConsumerApplicationShutdown shutdown = new ConsumerApplicationShutdown(applicationContext); + + shutdown.request(); + shutdown.request(); + + verify(applicationContext, times(1)).close(); + } +} diff --git a/backend/coinflow-consumer-app/src/test/java/com/coinflow/config/RedisConsumerGroupManagerTest.java b/backend/coinflow-consumer-app/src/test/java/com/coinflow/config/RedisConsumerGroupManagerTest.java new file mode 100644 index 00000000..46bf6524 --- /dev/null +++ b/backend/coinflow-consumer-app/src/test/java/com/coinflow/config/RedisConsumerGroupManagerTest.java @@ -0,0 +1,124 @@ +package com.coinflow.config; + +import com.coinflow.config.properties.TickConsumerProperties; +import java.nio.charset.StandardCharsets; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.data.redis.connection.RedisConnection; +import org.springframework.data.redis.connection.RedisStreamCommands; +import org.springframework.data.redis.connection.stream.ReadOffset; +import org.springframework.data.redis.core.RedisCallback; +import org.springframework.data.redis.core.RedisTemplate; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class RedisConsumerGroupManagerTest { + + @Mock + private RedisTemplate redisTemplate; + + @Mock + private RedisConnection connection; + + @Mock + private RedisStreamCommands streamCommands; + + @Mock + private ConsumerApplicationShutdown applicationShutdown; + + private RedisConsumerGroupManager groupManager; + + @BeforeEach + void setUp() { + TickConsumerProperties properties = new TickConsumerProperties( + "tick:raw", "tick-consumer-group", "consumer-1", 200_000L, 0.8); + groupManager = new RedisConsumerGroupManager(redisTemplate, properties, applicationShutdown); + } + + @Test + void createsGroupAndStreamAtomically() { + stubRedisExecute(); + + groupManager.ensureConsumerGroup(); + + ArgumentCaptor streamKey = ArgumentCaptor.forClass(byte[].class); + ArgumentCaptor readOffset = ArgumentCaptor.forClass(ReadOffset.class); + verify(streamCommands).xGroupCreate( + streamKey.capture(), + eq("tick-consumer-group"), + readOffset.capture(), + eq(true)); + + assertThat(new String(streamKey.getValue(), StandardCharsets.UTF_8)).isEqualTo("tick:raw"); + assertThat(readOffset.getValue().getOffset()).isEqualTo("$"); + } + + @Test + void acceptsExistingConsumerGroup() { + stubRedisExecute(); + when(streamCommands.xGroupCreate(any(byte[].class), any(), any(ReadOffset.class), eq(true))) + .thenThrow(new RuntimeException("BUSYGROUP Consumer Group name already exists")); + + assertThatCode(groupManager::ensureConsumerGroup).doesNotThrowAnyException(); + } + + @Test + void failsStartupForUnexpectedGroupInitializationError() { + stubRedisExecute(); + when(streamCommands.xGroupCreate(any(byte[].class), any(), any(ReadOffset.class), eq(true))) + .thenThrow(new RuntimeException("Redis connection failed")); + + assertThatThrownBy(groupManager::ensureConsumerGroup) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("Failed to initialize Redis consumer group"); + } + + @Test + void keepsSubscriptionActiveForNestedNoGroupError() { + RuntimeException error = new RuntimeException( + "Error in execution", new RuntimeException("NOGROUP No such consumer group")); + + assertThat(groupManager.shouldCancelSubscription(error)).isFalse(); + assertThat(groupManager.shouldCancelSubscription(new RuntimeException("timeout"))).isTrue(); + } + + @Test + void recreatesMissingGroupWithoutCancellingSubscription() { + stubRedisExecute(); + RuntimeException error = new RuntimeException("NOGROUP No such consumer group"); + + groupManager.handleSubscriptionError(error); + + verify(streamCommands).xGroupCreate(any(byte[].class), eq("tick-consumer-group"), any(ReadOffset.class), eq(true)); + } + + @Test + void shutsDownApplicationForFatalSubscriptionError() { + RuntimeException error = new RuntimeException("Redis command timeout"); + + groupManager.handleSubscriptionError(error); + + verify(applicationShutdown).request(); + } + + @SuppressWarnings("unchecked") + private void stubRedisExecute() { + when(connection.streamCommands()).thenReturn(streamCommands); + doAnswer(invocation -> { + RedisCallback callback = invocation.getArgument(0); + return callback.doInRedis(connection); + }).when(redisTemplate).execute(any(RedisCallback.class)); + } +} diff --git a/backend/coinflow-consumer-app/src/test/java/com/coinflow/monitoring/StreamLagMonitorWorkerTest.java b/backend/coinflow-consumer-app/src/test/java/com/coinflow/monitoring/StreamLagMonitorWorkerTest.java new file mode 100644 index 00000000..244fa24e --- /dev/null +++ b/backend/coinflow-consumer-app/src/test/java/com/coinflow/monitoring/StreamLagMonitorWorkerTest.java @@ -0,0 +1,77 @@ +package com.coinflow.monitoring; + +import com.coinflow.config.properties.TickConsumerProperties; +import java.util.concurrent.atomic.AtomicReference; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.data.redis.core.RedisTemplate; + +import static com.coinflow.monitoring.constant.MetricConstants.STREAM_BACKLOG_COUNT; +import static com.coinflow.monitoring.constant.MetricConstants.STREAM_BACKLOG_RETENTION_RATIO; +import static com.coinflow.monitoring.constant.MetricConstants.STREAM_RETENTION_WARNING_COUNT; +import static com.coinflow.monitoring.constant.MetricConstants.TAG_MODULE; +import static com.coinflow.monitoring.constant.MetricConstants.VALUE_MODULE_CONSUMER; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class StreamLagMonitorWorkerTest { + + @Mock + private RedisTemplate redisTemplate; + + @Mock + private MetricRecorder metricRecorder; + + private final AtomicReference backlogGauge = new AtomicReference<>(0.0); + private final AtomicReference retentionRatioGauge = new AtomicReference<>(0.0); + + private StreamLagMonitorWorker worker; + + @BeforeEach + void setUp() { + TickConsumerProperties properties = new TickConsumerProperties( + "tick:raw", "tick-consumer-group", "consumer-1", 200_000L, 0.8); + when(metricRecorder.registerGauge( + STREAM_BACKLOG_COUNT, 0.0, TAG_MODULE, VALUE_MODULE_CONSUMER)) + .thenReturn(backlogGauge); + when(metricRecorder.registerGauge( + STREAM_BACKLOG_RETENTION_RATIO, 0.0, TAG_MODULE, VALUE_MODULE_CONSUMER)) + .thenReturn(retentionRatioGauge); + + worker = new StreamLagMonitorWorker(redisTemplate, properties, metricRecorder); + worker.init(); + } + + @Test + void recordsBacklogRetentionRatio() { + worker.updateBacklogMetrics(100_000); + + assertThat(backlogGauge.get()).isEqualTo(100_000.0); + assertThat(retentionRatioGauge.get()).isEqualTo(0.5); + } + + @Test + void recordsWarningOnlyWhenThresholdIsCrossed() { + worker.updateBacklogMetrics(160_000); + worker.updateBacklogMetrics(180_000); + + verify(metricRecorder, times(1)).increment( + STREAM_RETENTION_WARNING_COUNT, + TAG_MODULE, + VALUE_MODULE_CONSUMER); + + worker.updateBacklogMetrics(100_000); + worker.updateBacklogMetrics(170_000); + + verify(metricRecorder, times(2)).increment( + STREAM_RETENTION_WARNING_COUNT, + TAG_MODULE, + VALUE_MODULE_CONSUMER); + } +} diff --git a/backend/coinflow-consumer-app/src/test/resources/application-test.yml b/backend/coinflow-consumer-app/src/test/resources/application-test.yml index 8ee99cbb..b6aff5d0 100644 --- a/backend/coinflow-consumer-app/src/test/resources/application-test.yml +++ b/backend/coinflow-consumer-app/src/test/resources/application-test.yml @@ -31,6 +31,9 @@ coinflow: redis: stream: tick: + enabled: false stream-key: test-stream group: test-group consumer-name: test-consumer + max-length: 200000 + lag-warning-ratio: 0.8 diff --git a/backend/coinflow-infra-redis/build.gradle b/backend/coinflow-infra-redis/build.gradle index 63d052a0..c25653ff 100644 --- a/backend/coinflow-infra-redis/build.gradle +++ b/backend/coinflow-infra-redis/build.gradle @@ -10,6 +10,8 @@ dependencies { implementation 'org.springframework.boot:spring-boot-starter-data-redis' implementation 'com.fasterxml.jackson.core:jackson-databind' + + testImplementation 'org.springframework.boot:spring-boot-starter-test' } bootJar { enabled = false } diff --git a/backend/coinflow-infra-redis/src/main/java/com/coinflow/publish/config/TickPublisherConfig.java b/backend/coinflow-infra-redis/src/main/java/com/coinflow/publish/config/TickPublisherConfig.java index e6123231..759033d4 100644 --- a/backend/coinflow-infra-redis/src/main/java/com/coinflow/publish/config/TickPublisherConfig.java +++ b/backend/coinflow-infra-redis/src/main/java/com/coinflow/publish/config/TickPublisherConfig.java @@ -3,9 +3,11 @@ import com.coinflow.tick.publisher.TickPublisher; import com.coinflow.publish.stream.RedisStreamTickPublisher; import com.coinflow.monitoring.MetricRecorder; +import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.util.Assert; @Configuration public class TickPublisherConfig { @@ -13,8 +15,12 @@ public class TickPublisherConfig { @Bean public TickPublisher tickPublisher( RedisTemplate rawRedisTemplate, - MetricRecorder metricRecorder + MetricRecorder metricRecorder, + @Value("${redis.stream.tick.stream-key:tick:raw}") String streamKey, + @Value("${redis.stream.tick.max-length:200000}") long maxLength ) { - return new RedisStreamTickPublisher(rawRedisTemplate, metricRecorder); + Assert.hasText(streamKey, "redis.stream.tick.stream-key must not be blank"); + Assert.isTrue(maxLength > 0, "redis.stream.tick.max-length must be greater than zero"); + return new RedisStreamTickPublisher(rawRedisTemplate, metricRecorder, streamKey, maxLength); } } diff --git a/backend/coinflow-infra-redis/src/main/java/com/coinflow/publish/stream/RedisStreamTickPublisher.java b/backend/coinflow-infra-redis/src/main/java/com/coinflow/publish/stream/RedisStreamTickPublisher.java index 0167880c..6ff700be 100644 --- a/backend/coinflow-infra-redis/src/main/java/com/coinflow/publish/stream/RedisStreamTickPublisher.java +++ b/backend/coinflow-infra-redis/src/main/java/com/coinflow/publish/stream/RedisStreamTickPublisher.java @@ -1,8 +1,6 @@ package com.coinflow.publish.stream; import com.coinflow.monitoring.MetricRecorder; -import static com.coinflow.monitoring.constant.MetricConstants.STREAM_PUBLISH_LATENCY; - import com.coinflow.publish.exception.PublishErrorCode; import com.coinflow.publish.exception.PublishException; import com.coinflow.tick.publisher.TickPublisher; @@ -16,7 +14,10 @@ import org.springframework.data.redis.connection.stream.StreamRecords; import org.springframework.data.redis.core.RedisTemplate; -import static com.coinflow.monitoring.constant.MetricConstants.*; +import static com.coinflow.monitoring.constant.MetricConstants.STREAM_PUBLISH_FAILURE_COUNT; +import static com.coinflow.monitoring.constant.MetricConstants.STREAM_PUBLISH_LATENCY; +import static com.coinflow.monitoring.constant.MetricConstants.TAG_MODULE; +import static com.coinflow.monitoring.constant.MetricConstants.VALUE_MODULE_COLLECTOR; /** * Redis Stream을 통해 바이너리 틱 데이터를 전송하는 구현체입니다. @@ -25,11 +26,12 @@ @RequiredArgsConstructor public class RedisStreamTickPublisher implements TickPublisher { - public static final String RAW_TICK_STREAM = "tick:raw"; public static final String RAW_PAYLOAD_FIELD = "p"; private final RedisTemplate rawRedisTemplate; private final MetricRecorder metricRecorder; + private final String streamKey; + private final long maxLength; /** * 최적화된 바이너리 방식 (Zero-POJO) @@ -38,16 +40,15 @@ public class RedisStreamTickPublisher implements TickPublisher { public void publish(byte[] rawData) { // MapRecord 생성 (String, String, byte[]) MapRecord record = StreamRecords.newRecord() - .in(RAW_TICK_STREAM) + .in(streamKey) .ofMap(Map.of(RAW_PAYLOAD_FIELD, rawData)); - // MAXLEN ~ 1,000,000 설정을 통한 자동 트리밍 (XAddOptions) - XAddOptions options = XAddOptions.maxlen(STREAM_MAX_LEN).approximateTrimming(true); + XAddOptions options = XAddOptions.maxlen(maxLength).approximateTrimming(true); RecordId recordId = executePublish(() -> rawRedisTemplate.opsForStream().add(record, options)); log.debug("Published raw tick data. stream={}, recordId={}, maxlen={}", - RAW_TICK_STREAM, recordId.getValue(), STREAM_MAX_LEN); + streamKey, recordId.getValue(), maxLength); } /** @@ -62,7 +63,11 @@ private RecordId executePublish(Callable publishAction) { } return recordId; } catch (Exception e) { - log.error("Failed to publish tick data to Redis Stream: {}", e.getMessage()); + metricRecorder.increment( + STREAM_PUBLISH_FAILURE_COUNT, + TAG_MODULE, + VALUE_MODULE_COLLECTOR); + log.error("Failed to publish tick data to Redis Stream", e); if (e instanceof PublishException) { throw (PublishException) e; } diff --git a/backend/coinflow-infra-redis/src/test/java/com/coinflow/publish/stream/RedisStreamTickPublisherTest.java b/backend/coinflow-infra-redis/src/test/java/com/coinflow/publish/stream/RedisStreamTickPublisherTest.java new file mode 100644 index 00000000..f4e9ea30 --- /dev/null +++ b/backend/coinflow-infra-redis/src/test/java/com/coinflow/publish/stream/RedisStreamTickPublisherTest.java @@ -0,0 +1,90 @@ +package com.coinflow.publish.stream; + +import com.coinflow.monitoring.MetricRecorder; +import java.util.concurrent.Callable; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.data.redis.connection.RedisStreamCommands.XAddOptions; +import org.springframework.data.redis.connection.stream.MapRecord; +import org.springframework.data.redis.connection.stream.RecordId; +import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.data.redis.core.StreamOperations; + +import static com.coinflow.monitoring.constant.MetricConstants.STREAM_PUBLISH_LATENCY; +import static com.coinflow.monitoring.constant.MetricConstants.STREAM_PUBLISH_FAILURE_COUNT; +import static com.coinflow.monitoring.constant.MetricConstants.TAG_MODULE; +import static com.coinflow.monitoring.constant.MetricConstants.VALUE_MODULE_COLLECTOR; +import static com.coinflow.publish.stream.RedisStreamTickPublisher.RAW_PAYLOAD_FIELD; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class RedisStreamTickPublisherTest { + + @Mock + private RedisTemplate redisTemplate; + + @Mock + private StreamOperations streamOperations; + + @Mock + private MetricRecorder metricRecorder; + + private RedisStreamTickPublisher publisher; + + @BeforeEach + void setUp() throws Exception { + publisher = new RedisStreamTickPublisher( + redisTemplate, metricRecorder, "tick:raw", 200_000L); + + when(redisTemplate.opsForStream()).thenReturn(streamOperations); + when(metricRecorder.recordTime( + eq(STREAM_PUBLISH_LATENCY), any(Callable.class), any(String[].class))) + .thenAnswer(invocation -> { + Callable callable = invocation.getArgument(1); + return callable.call(); + }); + } + + @Test + void publishesWithConfiguredApproximateMaxLength() { + byte[] payload = new byte[]{1, 2, 3}; + when(streamOperations.add(any(MapRecord.class), any(XAddOptions.class))) + .thenReturn(RecordId.of("1-0")); + + publisher.publish(payload); + + @SuppressWarnings("unchecked") + ArgumentCaptor> recordCaptor = + ArgumentCaptor.forClass(MapRecord.class); + ArgumentCaptor optionsCaptor = ArgumentCaptor.forClass(XAddOptions.class); + verify(streamOperations).add(recordCaptor.capture(), optionsCaptor.capture()); + + assertThat(recordCaptor.getValue().getStream()).isEqualTo("tick:raw"); + assertThat(recordCaptor.getValue().getValue().get(RAW_PAYLOAD_FIELD)).isEqualTo(payload); + assertThat(optionsCaptor.getValue().getMaxlen()).isEqualTo(200_000L); + assertThat(optionsCaptor.getValue().isApproximateTrimming()).isTrue(); + } + + @Test + void recordsPublishFailureWhenRedisRejectsXadd() { + when(streamOperations.add(any(MapRecord.class), any(XAddOptions.class))) + .thenThrow(new RuntimeException("OOM command not allowed")); + + assertThatThrownBy(() -> publisher.publish(new byte[]{1, 2, 3})) + .isInstanceOf(RuntimeException.class); + + verify(metricRecorder).increment( + STREAM_PUBLISH_FAILURE_COUNT, + TAG_MODULE, + VALUE_MODULE_COLLECTOR); + } +} diff --git a/infra/docker/.env.example b/infra/docker/.env.example new file mode 100644 index 00000000..19f3f85c --- /dev/null +++ b/infra/docker/.env.example @@ -0,0 +1,3 @@ +DOCKERHUB_USERNAME=your-dockerhub-username +DB_PASSWORD=replace-with-a-strong-database-password +REDIS_PASSWORD=replace-with-a-strong-redis-password diff --git a/infra/docker/docker-compose-prod.yml b/infra/docker/docker-compose-prod.yml index 3e161d81..fa288363 100644 --- a/infra/docker/docker-compose-prod.yml +++ b/infra/docker/docker-compose-prod.yml @@ -39,6 +39,7 @@ services: - SPRING_CONFIG_NAME=application-api - SPRING_PROFILES_ACTIVE=prod - SPRING_DATA_REDIS_HOST=redis + - SPRING_DATA_REDIS_PASSWORD=${REDIS_PASSWORD:?REDIS_PASSWORD is required} - SPRING_DATASOURCE_URL=jdbc:postgresql://postgres:5432/ohlc - SPRING_DATASOURCE_USERNAME=ohlc_user - SPRING_DATASOURCE_PASSWORD=${DB_PASSWORD:-1234} @@ -60,6 +61,7 @@ services: - SPRING_CONFIG_NAME=application-ws - SPRING_PROFILES_ACTIVE=prod - SPRING_DATA_REDIS_HOST=redis + - SPRING_DATA_REDIS_PASSWORD=${REDIS_PASSWORD:?REDIS_PASSWORD is required} entrypoint: [ "sh", "-c", "java -Dspring.config.name=application-ws -Dspring.profiles.active=prod -Xms128m -Xmx256m -jar app.jar" ] depends_on: redis: @@ -76,18 +78,27 @@ services: - SPRING_CONFIG_NAME=application-consumer - SPRING_PROFILES_ACTIVE=prod - SPRING_DATA_REDIS_HOST=redis + - SPRING_DATA_REDIS_PASSWORD=${REDIS_PASSWORD:?REDIS_PASSWORD is required} - SPRING_DATASOURCE_URL=jdbc:postgresql://postgres:5432/ohlc - SPRING_DATASOURCE_USERNAME=ohlc_user - SPRING_DATASOURCE_PASSWORD=${DB_PASSWORD:-1234} - REDIS_STREAM_TICK_STREAMKEY=tick:raw - REDIS_STREAM_TICK_GROUP=tick-consumer-group - REDIS_STREAM_TICK_CONSUMERNAME=consumer-1 + - REDIS_STREAM_TICK_MAXLENGTH=200000 + - REDIS_STREAM_TICK_LAGWARNINGRATIO=0.8 entrypoint: [ "sh", "-c", "java -Dspring.config.name=application-consumer -Dspring.profiles.active=prod -Xms256m -Xmx384m -jar app.jar" ] depends_on: postgres: condition: service_healthy redis: condition: service_healthy + healthcheck: + test: [ "CMD-SHELL", "wget -q -O /dev/null http://localhost:8081/actuator/health" ] + interval: 5s + timeout: 3s + retries: 20 + start_period: 30s restart: always collector-app: @@ -100,16 +111,19 @@ services: - SPRING_CONFIG_NAME=application-collector - SPRING_PROFILES_ACTIVE=prod - SPRING_DATA_REDIS_HOST=redis + - SPRING_DATA_REDIS_PASSWORD=${REDIS_PASSWORD:?REDIS_PASSWORD is required} - SPRING_DATASOURCE_URL=jdbc:postgresql://postgres:5432/ohlc - SPRING_DATASOURCE_USERNAME=ohlc_user - SPRING_DATASOURCE_PASSWORD=${DB_PASSWORD:-1234} - BINANCE_WEBSOCKET_BASEURL=wss://stream.binance.com:9443/stream - BINANCE_WEBSOCKET_SYMBOLS_0=btcusdt - BINANCE_WEBSOCKET_TRADESTREAMSUFFIX=@trade + - REDIS_STREAM_TICK_STREAMKEY=tick:raw + - REDIS_STREAM_TICK_MAXLENGTH=200000 entrypoint: [ "sh", "-c", "java -Dspring.config.name=application-collector -Dspring.profiles.active=prod -Xms128m -Xmx256m -jar app.jar" ] depends_on: consumer-app: - condition: service_started + condition: service_healthy redis: condition: service_healthy restart: always @@ -137,9 +151,28 @@ services: redis: image: redis:7.2 container_name: coinflow-redis - command: [ "redis-server", "--bind", "0.0.0.0", "--appendonly", "no", "--protected-mode", "no", "--maxmemory", "64mb", "--maxmemory-policy", "allkeys-lru" ] + environment: + REDIS_PASSWORD: ${REDIS_PASSWORD:?REDIS_PASSWORD is required} + command: + - redis-server + - --bind + - 0.0.0.0 + - --appendonly + - "yes" + - --appendfsync + - everysec + - --protected-mode + - "no" + - --requirepass + - ${REDIS_PASSWORD:?REDIS_PASSWORD is required} + - --maxmemory + - 64mb + - --maxmemory-policy + - noeviction + volumes: + - redis_data:/data healthcheck: - test: [ "CMD", "redis-cli", "ping" ] + test: [ "CMD-SHELL", "REDISCLI_AUTH=$$REDIS_PASSWORD redis-cli ping | grep -q PONG" ] interval: 3s timeout: 3s retries: 10 @@ -148,3 +181,4 @@ services: volumes: postgres_data: + redis_data: