pub struct KafkaProducer {
producer: FutureProducer,
}
pub struct KafkaConsumer {
consumer: StreamConsumer,
}
So what: Producer와 Consumer를 별도 구조체로 분리했다.
So why: Producer는 seoul.rs에서 데이터를 발행하고, Consumer는 main.rs에서 백그라운드 태스크로 Valkey에 저장한다. 둘의 lifespan이 다르기에 하나로 합치면 AppState에서 의존성이 생기는 현상이 일어나므로 분리해줬다.
pub async fn send_to_topic(
&self,
topic: &str,
key: &str,
payload: &str,
) -> Result<(), String>
So what: topic, key, payload 세 가지만 받는 단일 메서드로 인터페이스를 단순화했다
So why: send_to_topic 하나로 합치면 seoul.rs에서 토픽 이름만 바꿔서 재사용할 수 있다. Kafka 토픽이 늘어날 경우에도 kafka.rs를 건드리지 않아도 된다.
let valkey_key = format!("realtime:{}:{}", topic, key);
if let Err(e) = cache.set_realtime(&valkey_key, &payload).await {
eprintln!("Valkey 저장 실패 [{}]: {}", valkey_key, e);
} else {
println!("Kafka 수신 → Valkey 저장: {}", valkey_key);
}
So what: Kafka 토픽과 키를 조합해서 realtime:{topic}:{key} 형식으로 Valkey에 저장한다
So why: Key만 보고 어떤 데이터인지 바로 알 수 있다. Producer가 저장한 키와 Consumer가 읽는 키가 자동으로 일치하도록 했고 네임스페이스를 realtime:으로 시작해서 경로 캐시 키와 충돌하지 않는다.
None => {
println!("Kafka 스트림 종료");
break;
}None => {
println!("Kafka 스트림 종료");
break;
}
So what: 스트림이 종료되면 루프를 빠져나온다
So why: StreamConsumer의 스트림은 정상 운영 중에는 None이 나오지 않고 브로커 연결이 완전히 끊어지는 극단적인 상황에서 발생한다. 무한 루프로 두면 이 경우 CPU를 100% 사용하게 되며 스핀락 상태가 될 수 있어서 명시적으로 종료해야한다.
'Project > SSAFY2학기 특화 PJT' 카테고리의 다른 글
| [BE] fare.rs ( 요금 정보 ) (0) | 2026.03.13 |
|---|---|
| [BE] cache.rs ( Valkey 연동 ) (0) | 2026.03.13 |
| [BE] congestion.rs ( 혼잡도 AI ) (0) | 2026.03.13 |
| [BE] routes.rs ( TMAP API ) (0) | 2026.03.13 |
| [BE] seoul.rs ( 서울 열린 데이터 광장 API ) (0) | 2026.03.13 |
