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% 사용하게 되며 스핀락 상태가 될 수 있어서 명시적으로 종료해야한다.

+ Recent posts