요구 사항
- kafka로 들어오는 실시간 gps 데이터에 대한 실시간 데이터 처리 후 bigquery(data warehouse) 내 적재해야 함
- 분당 수백만 데이터에 대해 LAG 없이 처리되어야 함
- 데이터가 avro로 serialize되어서 들어오기 때문에 kafka schema registry로부터 schema 정보를 읽어 avro deserialize 수행해야 함
- 데이터 전송 간 network bandwidth resource 절감 및 전송 속도 개선을 위해 avro schema를 사용해 데이터를 serialize해 전송
- avro schema가 변경되어도 원활하게 데이터를 처리할 수 있어야 함
- 새로운 컬럼에 대해 backfill이 용이해야 함
개발 사항

- golang server를 사용한 consumer group 구현
- 실시간 속도 및 서버 리소스 절감을 위해 golang server로 구현
- golang의 sarama library에는 deserializer 설정과 같은 기능이 없기 때문에 직접 deserializer를 구현
- golang 서버는 Deserialize, Data Process, Publish 부분으로 나뉘어져 있음
- Deserialize : Avro 데이터를 Deserialize하는 부분
- Data Process : gps 데이터를 BigQuery에 적재하기 용이한 형태로 변환하는 부분
- Publish : Pub/Sub으로 Publishing하는 부분
- 실제 데이터 적재는 bigquery subscribe를 사용해 구현
- Limiter를 사용해 resource exhausting을 방지
Deserializer
kafka topic으로부터 consume한 데이터를 avro schema에 맞춰 deserialize
- 데이터의 schema가 변경되어도 원활하게 데이터가 적재될 수 있도록 함
- 초기 Deserializer 생성 시, kafka schema registry로부터 가장 최신의 avro schema 정보를 메모리에 로드
- consume한 데이터의 avro schema version을 확인해 메모리에 해당 version의 avro schema가 존재하는지 확인
- avro schema가 존재하지 않는다면 kafka schema registry로부터 schema 정보를 메모리에 로드
- 읽어온 avro schema를 기반으로 데이터 deserialize 수행
func (deser *AvroDeserializer) Deserialize(data []byte) (interface{}, error) {
// 01. consume한 데이터의 avro schema version을 확인해 메모리에 해당 version의 avro schema가 존재하는지 확인
var schemaID int32
err := binary.Read(bytes.NewReader(data[1:5]), binary.BigEndian, &schemaID)
if err != nil {
return nil, errors.Wrapf(err, "failed to decode schemaID. schemaID : %d\\n", schemaID)
}
// 02. 해당 데이터의 avro schema version에 맞는 avroCodec 선택
// 만약 해당 version의 avroCodec이 존재하지 않는다면 kafka schema registry로부터 schema 정보를 메모리에 로드
codec, err := deser.GetAvroCodecBySchemaID(int(schemaID))
if err != nil {
err = errors.Wrapf(err, "failed to get codec by schemaID. schemaID : %d\\n", schemaID)
return nil, err
}
// 03. avroCodec을 통해 data deserialize 수행
deserializedData, _, err := codec.NativeFromBinary(data[5:])
if err != nil {
err = errors.Wrap(err, "failed to decode data.")
return nil, err
}
return deserializedData, nil
}
Data Process
gps 데이터를 적절한 형태의 데이터로 처리
- 각 linkID 별로 데이터를 flattening
- 추후 backfill이 용이하도록 raw 데이터를 같이 적재
type FlattenGps struct {
...
// backfill을 위한 해당 path의 raw 데이터 저장
Path string `json:"path"`
}
Publisher
gps 데이터를 pub/sub으로 push
func (p *GcpPubSubPublisher) Publish(publishData interface{}) (err error) {
// 01. marshaling publishData
jsonData, err := json.Marshal(publishData)
if err != nil {
err = errors.Wrap(err, "failed to marshal publishData.")
return
}
msg := &Message{
Data: jsonData,
}
// 02. data push and get result
for retry := 0; retry < p.Retry; retry++ {
if _, err = p.PublishTopic.Publish(p.PublishCtx, msg).Get(p.PublishCtx); err == nil {
break
}
}
if err != nil {
err = errors.Wrap(err, "failed to publish.")
return
}
return
}
ConsumerGroupHandler
ConsumerGroup에서 consume한 message를 처리하기 위한 handler
func (h *NaviGpsProcessorConsumerGroupHandler) ConsumeClaim(session ConsumerGroupSession, claim ConsumerGroupClaim) error {
for message := range claim.Messages() {
// 01. Acquire Resource
h.consumerGroupLimiter <- struct{}{}
go func(message *ConsumerMessage) {
session.MarkMessage(message, "")
// 02. Release Resource
defer func() { <-h.consumerGroupLimiter }()
// 03. Avro Deserialize
deserializedValue, err := h.Deserializer.Deserialize(message.Value)
if err != nil {
log.WithError(err).Error("message deserialize에 실패하였습니다.")
return
}
// 04. Flattening
flattenGpsList, err := dto.ToFlattenGps(deserializedValue)
if err != nil {
log.WithError(err).Error("Parsing Error가 발생하였습니다.")
return
}
// 05. Publishing
for _, flatGps := range flattenGpsList {
err = h.Publisher.Publish(flatGps)
if err != nil {
log.WithError(err).Errorf("Publishing에 실패하였습니다.")
return
}
}
}(message)
}
return nil
}
고려 사항
메시지 전달 보장 수준(Message Delivery Semantics)
데이터의 정합성이 크게 중요하지 않으므로 엄격한 보증 수준을 고려하지 않음
하지만 해당 데이터의 경우 실시간성과 처리 속도가 중요하기 때문에 at least once 수준에 가깝도록 설계
- 최대한 재처리를 하지 않도록 해 메세지 소비를 빠르게 하기 위함
config.Consumer.Offsets.AutoCommit.Enable = true
config.Consumer.Offsets.AutoCommit.Interval = time.Millisecond * 100
- ConsumerGroupHandler에서 데이터 처리 전에 MarkMessage 처리
- autoCommit interval의 주기를 짧게 설정
Publisher Publish / Get 동시 처리 여부
as-is : Publisher에서 데이터 Publish와 Get을 동시에 처리하도록 설정
- p.PublishTopic.Publish(p.PublishCtx, msg).Get(p.PublishCtx)
Get(p.PublishCtx)의 경우 publish가 완료될 때까지 해당 goroutine이 blocking되므로 실시간성이 떨어진다.
하지만 Publish와 Get을 분리한다면 publish 실패 시 재처리가 어려워진다.
- 따라서 publish 시 retry 로직 또는 별도의 error handling을 통한 재처리가 중요하다면 Publish와 Get을 동시에 처리
- 실시간성이 중요하다면 Publish와 Get의 로직을 분리하는 것도 고려
해당 어플리케이션은 실시간성이 중요하지만 publish와 Get을 분리하였을 때의 코드 복잡성과 실시간성의 이점을 고려하였을 때 publish와 Get을 분리하지 않도록 설정
- consuming 속도보다 publishing 속도가 현저히 느린 경우, publishing queue에 LAG이 발생하면서 timeout이 발생하는 이슈도 존재
ConsumerGroup 병렬처리
1. ConsumerGroup 내 단일 처리
- 성능 : 200/minutes
- 원인 : Publisher 내 Get(publishCtx) 로 인한 goroutine blocking
_, err := publishTopic.Publish(publishCtx, msg).Get(publishCtx)
2. Publisher 내 Publish/Get 로직 분리
- 성능 : 100k/minutes
publishResult := publishTopic.Publish(publishCtx, msg)
publishResultChan <- publishResult
publishing 시 lag이 생기면 publishing error가 발생하는 이슈 존재
publishing에 실패하였습니다. error : context deadline exceeded
- consuming 속도보다 publishing 속도가 현저하게 낮을 경우, grpc channel에 lag이 발생하면서 처리 지연이 발생
- 처리 지연이 설정한 timeout보다 높아지는 경우 context가 종료되며 에러가 발생
- grpc channel은 server의 threads 수 증가와 관련있기 때문에 무한정 늘릴 수 없다.
-> 결국은 consuming 속도를 제어해야 함
3. Goroutine Parallel Processing
- ConsumerGroup 내에서 goroutine으로 consuming을 병렬처리
- publish/Get 로직은 통합하여 각 goroutine마다 blocking될 수는 있지만 병렬처리로 인한 throughput 향상
Goroutine 제어를 하지 않는 경우 Resource Exhausting 문제 발생
- Channel을 활용해 Goroutine pool을 관리하는 Limiter 추가
- session 초기화 및 제거 시에 limiter 생성 및 삭제
func (h *NaviGpsProcessorConsumerGroupHandler) Setup(_ ConsumerGroupSession) error {
h.consumerGroupLimiter = make(chan struct{}, *configs.ConsumerGroupParallelism)
return nil
}
func (h *NaviGpsProcessorConsumerGroupHandler) Cleanup(_ ConsumerGroupSession) error {
for i := 0; i < cap(h.consumerGroupLimiter); i++ {
h.consumerGroupLimiter <- struct{}{}
}
close(h.consumerGroupLimiter)
return nil
}
goroutine의 개수는 할당된 resource에 맞게 적절한 설정이 필요
- goroutine의 개수가 증가할수록 cpu 및 memory 사용량이 높아진다.
4. TO BE : Goroutine Parallel Processing + Publish/Get 로직 분리
- 분리 시에 성능 상 많은 이점이 있지만 publishing 속도 및 grpc channel 설정을 consuming 속도에 맞춰 조정해나가야 함
- 이러한 방법은 개발 리소스가 많이 들고, 서버 안정성이 저하됨
- 이전 방법으로도 상당히 높은 성능을 낼 수 있기 때문에 해당 방법은 적용하지 않음.
- 추후 추가 성능 개선 필요 시 고려해 볼 수 있음
Consumer Configuration
- 실시간성을 위해 maxWaitTime을 100ms로 설정
- Fetch.Min의 경우 kafka broker와의 RTT를 줄이기 위해 기본 1024 batch size를 설정
- channel buffer size의 경우 메모리 상황에 맞게 설정
- maxWaitTime 및 fetch.min에 의해 설정된 주기로 데이터를 consume한 값을 channel buffer size에 담아둔다.
- 이때 channel buffer가 꽉 차게되면 consuming을 멈춘다.
Producer Configuration
- RTT를 줄이기 위해 DelayThreshold, CountThreshold, ByteThreshold 값 조정
- FlowControlSettings는 모두 제거, 별도의 monitoring(prometheus + grafana)를 통한 에러 감지를 수행
- FlowControl로 인해 publishing이 안 되는 것보다 LAG이 발생하는 것이 서버 안정성 측면에서 좋다고 판단
- bufferedByteLimit의 경우 메모리 상황에 맞게 설정
결과
- cpu 4, memory 4Gi node 한 대당 분당 몇십만 데이터 처리
- avro schema 변경에 대응하기 위해 deserializer 내에서 kafka schema registry에서 schema 정보를 조회해 메모리에 적재
- 신규 컬럼에 대한 backfill을 용이하게 하기 위해 linkID 당 raw 데이터를 bigquery에 같이 적재
'미니프로젝트' 카테고리의 다른 글
| 실시간 게시글 키워드 기반 푸시 알람 어플리케이션 설계 (1) | 2024.01.17 |
|---|