Spring Data Reids에서 Redis Streams를 구현
- Redis Streams는 redis pub/sub과 다르게 message queueing system
- 즉 queue 형태의 자료구조에 publisher가 append하면, consumer가 데이터를 읽는다.
- 이때 consumer group 단위로 데이터를 읽을 수 있다.
- redis pub/sub과 다르게 redis 내의 메모리를 사용해 데이터를 저장하기 때문에 consumer는 놓친 부분부터 데이터를 다시 읽을 수 있다.
publishing
spring을 통해서 redis streams에 데이터를 append할 수 있다.
- low-level의 RedisConnection을 사용할 수도 있고,
- high-level의 StreamOperations을 사용할 수도 있다.
여기에서는 StreamOperations을 사용한 방식만 소개한다.
StringRecord 사용
@Test
public void publishingTest() throws JsonProcessingException {
// given
RedisExampleEntity entity = RedisExampleEntity.builder()
.id(1L)
.callId(1L)
.status("OK")
.build();
String entityString = objectMapper.writeValueAsString(entity);
// when
StringRecord record = StreamRecords.newRecord().ofStrings(Map.of("redis-example-entity", entityString)).withStreamKey("test-stream");
RecordId id = redisTemplate.opsForStream().add(record);
// then
System.out.println(id);
}
output
[
{
"id": "1705075266513-0",
"redis-example-entity": "{\\"id\\":1,\\"callId\\":1,\\"status\\":\\"OK\\"}"
}
]
ObjectRecord 사용
@Test
public void publishingTest() throws JsonProcessingException {
// given
RedisExampleEntity entity = RedisExampleEntity.builder()
.id(1L)
.callId(1L)
.status("OK")
.build();
// when
ObjectRecord<String, RedisExampleEntity> record = StreamRecords.newRecord().ofObject(entity).withStreamKey("test-stream");
RecordId id = redisTemplate.opsForStream().add(record);
// then
System.out.println(id);
}
output
[
{
"id": "1",
"_class": "com.moon.redisexample.entity.RedisExampleEntity",
"callId": "1",
"status": "OK"
}
]
- 클래스 정보가 같이 저장된다.
StringRecord vs ObjectRecord
StringRecord
- 직접 serialize, deserialize를 수행해줘야 한다.
- 클래스 정보가 노출되지 않고, subscriber가 클래스 정보를 몰라도 된다.
- 하나의 레코드에 다양한 이벤트 값을 한번에 publishing할 수 있다.
ObjectRecord
- 별도의 serialize, deserialize를 수행할 필요가 없다.
- 클래스 정보가 노출되고, subscriber가 클래스 정보를 알아야 한다.
- 클래스의 위치가 바뀌면 deserialize 시 문제가 발생할 수 있다.
- 하나의 레코드는 하나의 이벤트만 publishing할 수 있다.
Subscribing
spring을 통해서 실시간으로 redis stream의 데이터를 consuming/subscribing 할 수 있다.
이 역시 low-level의 RedisConnection, high-level의 StreamOperations 둘 다 존재
RedisConnection과 StreamOperations 모두 동기 방식으로, 해당 thread를 block 시킬 수 있다.
- read timeout이 발생하거나, 메세지를 받게 되면 block이 풀린다.
이러한 blocking 방식은 서비스의 안정성을 떨어트리게 된다.
- 따라서 Redis Pub/Sub과 동일하게 subscribing 시 별도의 쓰레드풀을 관리해 동작할 수 있도록 해주는 StreamListenerContainer를 제공
StreamListenerContainer
StreamListenerContainer는 크게 다음 3가지 설정이 필요
- RedisConnectionFactory : 데이터를 전달할 RedisConnectionFactory
- StreamMessageListenerContainerOptions : StreamListenerContainer에서 사용할 설정값
- ex) batchSize, pollTimeout, …
- StreamListener : 실제 consuming 작업을 수행할 object
- 메세지를 처리하고, acknowledge 체크 등을 수행
@Bean
public StreamMessageListenerContainer<String, MapRecord<String, String, String>> listenerContainer(
@Autowired StreamListener<String, MapRecord<String, String, String>> streamListener
) {
// StreamListenerContainer의 옵션 설정
StreamMessageListenerContainerOptions<String, MapRecord<String, String, String>> containerOptions = StreamMessageListenerContainerOptions
.builder()
.pollTimeout(Duration.ofMillis(100))
.batchSize(1)
.build();
// StreamListenerContainer 생성
StreamMessageListenerContainer<String, MapRecord<String, String, String>> listenerContainer =
StreamMessageListenerContainer.create(connectionFactory, containerOptions);
// consumer(streamListener) 등록
listenerContainer.receive(
Consumer.from(groupId, consumerId),
StreamOffset.create(streamId, ReadOffset.lastConsumed()),
streamListener);
// subscribing 시작
listenerContainer.start();
return listenerContainer;
}
StreamListener
@RequiredArgsConstructor
class TestStreamListener implements StreamListener<String, MapRecord<String, String, String>> {
final RedisTemplate<String, String> redisTemplate;
final String groupId;
final String consumerId;
@Override
public void onMessage(MapRecord<String, String, String> message) {
System.out.println("message.id = " + message.getId());
System.out.println("message.stream = " + message.getStream());
System.out.println("message.value = " + message.getValue());
redisTemplate.opsForStream().acknowledge(groupId, message);
}
}
- StreamListenerContainer에서 receive로 등록하게 되면 autoAcknowledge가 false이므로, StreamListener에서 acknowledge 처리를 해줘야 한다.
- autoAcknowledge를 설정하려면 register를 통해 등록해 설정하면 된다.
'Spring Data Redis(스프링 데이터 레디스)' 카테고리의 다른 글
| Redis Pub/Sub (0) | 2023.12.29 |
|---|---|
| Redis Cache (0) | 2023.12.29 |
| Redis Serializers (1) | 2023.12.29 |
| RedisTemplate (0) | 2023.12.29 |
| Redis Connection Modes (0) | 2023.12.29 |