이번 구현 목표는 "외부 시스템으로부터 시세 정보를 웹소켓을 구독해서 메시지 브로커에 전달"하는 것이다.

스프링을 백엔드로 하는 플랫폼 웹서버와 분리가 되었기 때문에 다른 언어와 프레임워크로 시세 수신 서버를 구현을 해도 된다는 선택지가 생겼다.
그래서 직접 벤치마크 테스트를 해서 언어와 프레임워크를 선택할까? 했지만, 너무 산으로 가는 것 같았다.
보통 어느 언어와 프레임워크를 사용하는 지 자료를 찾다가 웹소켓 벤치마크 관련 논문자료를 찾았다.
여러 조건하에서 벤치마크 테스트를 한 내용이다.
Python은 커넥션이 많아질수록 높은 지연시간, C는 특정구간에서 안정적이지 못한 결과를 가져왔다.
결론으로 Python은 소켓통신으로 비추천이라고 한다.
서비스 환경에 맞게 조건들 조정하고, 자체적으로 테스트해보는 게 확실한데...
참고만하고 그냥 Java, Spring로 가자.
눈에 보이는 것들은 Redis, Kafka가 있는데, 토스 세미나 영상에서 메세지 브로커 처리 속도를 비교한 것이 있다.

출처: https://www.youtube.com/watch?v=SF7eqlL0mjw
Redis가 Kafka에 비해 낮은 지연시간, Queue가 없는 PUB/SUB구조라는 장점이 있는데,
대신 Redis로부터 데이터를 수신받는 서버가 처리를 제때 못한다면 메시지가 유실된다는 단점이 있다.
토스는 이걸 Spring의 Redis 라이브러리 구조를 분석하고 이벤트 루프를 구현해서 해결한 것 같은데...
Redis로 가자
필요한 것은 무언인지 생각을 해보자.
구현 요소
1. 서버 내에서 외부 시스템으로부터 데이터를 불러오기
2. 불러온 데이터에서 필요한 데이터만 가공
3. 가공된 데이터를 메시지 브로커로 Publish
이 3가지 기능만 제대로 작동하면 시세수신 서버 구현은 완료다.

Websocket 의존성만 추가해서 생성했다.
Redis 의존성도 추가해야하는데, 서버 내에서 외부 시스템으로부터 데이터를 불러오는 과정을 먼저 구현해보고 눈으로 보고 싶어서 Webscoket만 추가했다.
Redis 의존성은 있다가 추가한다.
https://bybit-exchange.github.io/docs/v5/websocket/public/ticker
실시간 시세정보를 받고 싶은 코인 티커를 토픽으로 서버에 전달하면 관련 정보를 웹소켓을 통해 수신할 수 있다.
모의투자플랫폼에서 지원하는 거래는 현물거래가 아닌 선물거래로 진행할 예정이라 "Linear/Inverse" 기준으로
ts, symbol, lastPrice 이 3가지를 가져올 것이다.
설명
1. ts: 시점 (timestamp)
2. symbol: 코인 이름 + 거래 코인단위
3. lastPrice: 마지막 가격
실제로 받아오는 데이터 로그를 찍어보면 시세변동정보뿐만이 아니라 실제 주문이 등록되는 정보들도 수신된다.
시세 정보가 수신되었을때만 Publish 하는 것이 필요하다.
package com.example.receivePrice.infra.websocket.bybit;
import jakarta.annotation.PostConstruct;
import org.springframework.stereotype.Component;
import org.springframework.web.socket.client.WebSocketClient;
import org.springframework.web.socket.client.standard.StandardWebSocketClient;
@Component
public class BybitWebSocketClient {
private static final String BYBIT_WS_URL = "wss://stream.bybit.com/v5/public/linear";
private final WebSocketClient webSocketClient = new StandardWebSocketClient();
private final BybitWebSocketHandler handler;
public BybitWebSocketClient(BybitWebSocketHandler handler) {
this.handler = handler;
}
@PostConstruct
public void connect() {
webSocketClient.execute(handler, BYBIT_WS_URL);
System.out.println("Connecting to Bybit WebSocket...");
}
}
외부 웹소켓을 연결하는 부분이다.
BybitWebSocketClient 클래스는 스프링에서 @Component로 선언되어 서버가 올라갈때, 스프링 컨테이너 내에 Bean으로 등록이 된다.
Bean 생성과 의존성 주입이 완료된 직후 @PostConstruct가 붙은 메서드가 실행되고, 이 시점에 외부 Bybit WebSocket 서버로의 연결이 자동으로 이루어진다.
생성자에서 BybitWebSocketHandler를 주입받는데, 이는 BybitWebSocketHandler가 Spring 컨테이너에 Bean으로 등록되어 있어 Spring이 미리 생성해 둔 Bean 인스턴스를 자동으로 주입해준다.
package com.example.receivePrice.infra.websocket.bybit;
import com.example.receivePrice.infra.redis.RedisPricePublisher;
import org.springframework.stereotype.Component;
import org.springframework.web.socket.TextMessage;
import org.springframework.web.socket.WebSocketSession;
import org.springframework.web.socket.handler.TextWebSocketHandler;
import tools.jackson.databind.JsonNode;
import tools.jackson.databind.ObjectMapper;
import java.util.List;
import java.util.Map;
@Component
public class BybitWebSocketHandler extends TextWebSocketHandler {
private final ObjectMapper objectMapper = new ObjectMapper();
private final RedisPricePublisher publisher;
public BybitWebSocketHandler(RedisPricePublisher publisher) {
this.publisher = publisher;
}
// 소켓 연결 이후 서버에 보낼 메세지
@Override
public void afterConnectionEstablished(WebSocketSession session) throws Exception {
// subscribe 메세지 (소켓 구독 방법: 문서 참조)
Map<String, Object> subscribeMsg = Map.of(
"op", "subscribe",
"args", List.of("tickers.BTCUSDT", "tickers.ETHUSDT")
);
// 연결된 소켓에 메세지 보내기
session.sendMessage(
new TextMessage(objectMapper.writeValueAsString(subscribeMsg))
);
System.out.println("Subscribed to BTCUSDT, ETHUSDT ticker");
}
// 수신 받은 메세지 가공
@Override
protected void handleTextMessage(WebSocketSession session, TextMessage message) {
try {
JsonNode root = objectMapper.readTree(message.getPayload());
// 0. ts 노드 확인
JsonNode timeNode = root.get("ts");
if (timeNode == null) return;
String receiveTime = timeNode.asString();
// 1. data 노드 없으면 무시
JsonNode dataNode = root.get("data");
if (dataNode == null) return;
// 2. symbol
JsonNode symbolNode = dataNode.get("symbol");
if (symbolNode == null) return;
String symbol = symbolNode.asString();
// 3. lastPrice
JsonNode lastPriceNode = dataNode.get("lastPrice");
if (lastPriceNode == null) return;
String lastPrice = lastPriceNode.asString();
System.out.println("symbol=" + symbol + ", price=" + lastPrice + ", time=" + receiveTime);
// Redis publish
publisher.publish(symbol, lastPrice, receiveTime);
} catch (Exception e) {
e.printStackTrace();
}
}
@Override
public void handleTransportError(WebSocketSession session, Throwable exception) {
System.err.println("WebSocket error: " + exception.getMessage());
}
}
TextWebSocketHandler를 상속받아와서 구현하였다.
https://docs.spring.io/spring-framework/docs/current/javadoc-api/org/springframework/web/socket/handler/TextWebSocketHandler.html
TextWebSocketHandler에서 상속받는 메서드들은 여기서 확인할 수 있다.
상속받는 메서드들 중 3가지 메서드를 재정의했다.
재정의 메서드 3가지
- afterConnectionEstablished
WebSocket 연결이 성공적으로 완료된 직후 구독 요청 메세지를 보낸다.
연결 session을 저장한다.- handleTextMessage
WebSocket으로 텍스트 메시지를 수신할 때마다 실시간으로 반복 호출되는 부분으로 수신 데이터를 처리한다.- handleTransportError
WebSocket 통신 중 에러 발생 시 처리하는 부분으로 session을 닫거나 연결을 재호출한다.
@Override
public void afterConnectionEstablished(WebSocketSession session) throws Exception {
// subscribe 메세지 (소켓 구독 방법: 문서 참조)
Map<String, Object> subscribeMsg = Map.of(
"op", "subscribe",
"args", List.of("tickers.BTCUSDT", "tickers.ETHUSDT")
);
// 연결된 소켓에 메세지 보내기
session.sendMessage(
new TextMessage(objectMapper.writeValueAsString(subscribeMsg))
);
System.out.println("Subscribed to BTCUSDT, ETHUSDT ticker");
}
웹소켓 연결이후 구독 정보를 외부 서버로 보내고, 그 연결 세션을 저장하는 부분이다.
Json 형태로 구독 정보를 보내야되는데, Map을 활용해 Bybit에서 요구하는 형태를 맞추었다.
// 수신 받은 메세지 가공
@Override
protected void handleTextMessage(WebSocketSession session, TextMessage message) {
try {
JsonNode root = objectMapper.readTree(message.getPayload());
// 0. ts 노드 확인
JsonNode timeNode = root.get("ts");
if (timeNode == null) return;
String receiveTime = timeNode.asString();
// 1. data 노드 없으면 무시
JsonNode dataNode = root.get("data");
if (dataNode == null) return;
// 2. symbol
JsonNode symbolNode = dataNode.get("symbol");
if (symbolNode == null) return;
String symbol = symbolNode.asString();
// 3. lastPrice
JsonNode lastPriceNode = dataNode.get("lastPrice");
if (lastPriceNode == null) return;
String lastPrice = lastPriceNode.asString();
System.out.println("symbol=" + symbol + ", price=" + lastPrice + ", time=" + receiveTime);
// Redis publish
publisher.publish(symbol, lastPrice, receiveTime);
} catch (Exception e) {
e.printStackTrace();
}
}
수신받은 데이터를 처리하는 부분이다.
Json형태의 데이터를 받아오기 때문에 JsonNode 객체형태로 전처리를 했고, 전처리된 데이터를 Redis로 Publish 했다.
publisher는 어디서 구현되었냐고? Redis 연동파트에서 다시 언급하겠다.
트러블슈팅
String symbol = symbolNode.toString();
-> String symbol = symbolNode.asString();
toString()으로 하니깐 값에 "" 큰따옴표가 붙어 Redis-cli 창에서 채널을 구독이 정상적으로 되지않았다.
JsonNode 객체의 값만 그대로 사용하고 싶으면 asText()로 변환한다고 한다.
API 정리를 위해 asString() 같은 새로운 이름이 도입되었고, asText()와 같은 역할을 하기때문에,
asString()으로 변환해주었다.
@Override
public void handleTransportError(WebSocketSession session, Throwable exception) {
System.err.println("WebSocket error: " + exception.getMessage());
}
일단 시세수신서버를 분리해놓기도 했고, 데이터가 메세지브로커까지 가는 과정을 확인해보는 것이 1순위라 연결 오류 정보만 띄울 수 있게 해놓았다. 재연결로직은 나중에...
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
프로젝트에 Redis 의존성을 추가해준다.
# Setting Redis Connection
spring.data.redis.host=${REDIS_HOST}
spring.data.redis.port=${REDIS_PORT}
spring.data.redis.password=${REDIS_PASSWORD}
Redis 서버 정보를 properties에 추가해준다.
Spring Boot의 Redis Auto Configuration이 해당 설정을 기반으로 RedisConnectionFactory를 생성한다.
이후 RedisConfig에서는 이 RedisConnectionFactory를 주입받아 Redis와의 연결을 수행하는 RedisTemplate를 구성한다.
properties → Auto Configuration → RedisConnectionFactory → RedisTemplate
RedisConfig가 뭔데? 아래있다.
package com.example.receivePrice.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.data.redis.serializer.StringRedisSerializer;
@Configuration
public class RedisConfig {
@Bean
public RedisTemplate<String, String> redisTemplate(RedisConnectionFactory connectionFactory) {
RedisTemplate<String, String> template = new RedisTemplate<>();
template.setConnectionFactory(connectionFactory);
template.setKeySerializer(new StringRedisSerializer());
template.setValueSerializer(new StringRedisSerializer());
return template;
}
}
Redis 서버와 통신할때 필요한 정보를 사전에 정의한 부분이다.
서버 정보, 데이터 저장, 전송 형태가 정의되어있다.
package com.example.receivePrice.infra.redis;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Component;
import tools.jackson.databind.ObjectMapper;
import java.util.Map;
@Component
public class RedisPricePublisher {
private final RedisTemplate<String, String> redisTemplate;
private final ObjectMapper objectMapper = new ObjectMapper();
public RedisPricePublisher(RedisTemplate<String, String> redisTemplate) {
this.redisTemplate = redisTemplate;
}
public void publish(String symbol, String price, String ts) throws Exception {
String channel = "price:" + symbol;
Map<String, String> payload = Map.of(
"symbol", symbol,
"price", price,
"ts", ts
);
redisTemplate.convertAndSend(
channel,
objectMapper.writeValueAsString(payload)
);
}
}
Redis Pub/Sub을 통해 시세 데이터를 발행하는 Publisher 부분이다.
코인별로 채널을 분리하고, JSON 형태의 메시지를 Map으로 구성한 후 RedisTemplate를 사용해 Redis로 전송한다.

비트코인, 이더리움 실시간 시세정보가 Redis PUB/SUB으로 올라가는 걸 확인할 수 있다.
Redis-cli에서 원하는 데이터가 올라올때 신기했다. 이게되네...
코드 작성은 GPT, 구글의 힘으로 어렵진 않았는데,
왜 이런것들을 쓰는지, 다른 방법은 없는지 납득을 하는 과정이 오래 걸렸다.
그리고 코드 내에 의존성 분리를 해야하는 부분들이 보인다.
일단 다 만들어놓고 이후 수정하는 방향으로 기술 부채를 쌓아두자.
지금부터 건들이면 완주가 안된다.
번외: 정리 금방 끝날 줄 알았는데...시간... 토익 언제함