본문으로 건너뛰기
홈
기술
기술 전체
프로그래밍68
컴퓨터 과학63
AI48
웹 개발36
인프라33
데이터31
소프트웨어 공학18
소개
← 목록으로데이터 › 메시징 › Kafka

Kafka REST Proxy

목차

Kafka REST Proxy

아래 명령은 2019년 Confluent REST Proxy v2 API 실습 기록이다. 현재 Confluent 문서에는 v2와 v3가 함께 있으며 API 경로·미디어 타입이 다르다. 배포한 REST Proxy 버전과 활성화된 API를 먼저 확인한다.

Kafka REST Proxy는 Confluent Platform 으로 HTTP를 통해 카프카와 연결하여 송 수신하는 역할을 한다.

준비사항

컨플루언트 플랫폼이 실행중 이어야한다.

$ confluent start kafka-rest
1. JSON 형식의 메시지로 해당 토픽에 전송
$ curl -X POST -H "Content-Type: application/vnd.kafka.json.v2+json" \
      -H "Accept: application/vnd.kafka.v2+json" \
      --data '{"records":[{"value":{"hello":"world"}}]}' \
      "http://localhost:8082/topics/source-topic"
2. JSON 데이터를 받기위해 컨슈머 생성
$ curl -X POST -H "Content-Type: application/vnd.kafka.v2+json" \
        --data '{"name": "my_json_consumer", "format": "json", "auto.offset.reset": "earliest"}' \
        http://localhost:8082/consumers/source-topic

생성 결과

{"instance_id":"my_json_consumer","base_uri":"http://kafka-rest-proxy:8082/consumers/source-topic/instances/my_json_consumer"}%
3. source-topic 토픽 my_json_consumer를 구독
$ curl -X POST -H "Content-Type: application/vnd.kafka.v2+json" --data '{"topics":["source-topic"]}' \
 http://localhost:8082/consumers/source-topic/instances/my_json_consumer/subscription
4. 데이터 요청
$ curl -X GET -H "Accept: application/vnd.kafka.json.v2+json" \
 http://localhost:8082/consumers/source-topic/instances/my_json_consumer/records
5. 컨슈머 종료
$ curl -X DELETE -H "Content-Type: application/vnd.kafka.v2+json" \
  http://localhost:8082/consumers/source-topic/instances/my_json_consumer           
6. 토픽의 목록을 가져온다
$ curl "http://localhost:8082/topics"
7. 하나의 토픽에 대한 정보를 가져온다
#ex) http://localhost:8082/topics/Topic-name

$ curl http://localhost:8082/topics/source-topic
8. 토픽 파티션 정보를 가져온다.
curl http://localhost:8082/topics/source-topic/partitions

자세한 사용법은 Confluent REST Proxy 에서 확인 할 수 있다.

HTTP 요청이 Kafka 레코드가 되기까지

REST Proxy는 HTTP 클라이언트와 Kafka 브로커 사이에서 프로듀서·컨슈머 API를 대신 호출한다. 브라우저나 다른 언어 서비스가 HTTP만 사용할 수 있는 환경에서는 편하지만, 프록시 배포·인증·연결 수·요청 크기·지연을 추가로 운영해야 한다. 네이티브 Kafka 클라이언트가 가능한 서비스라면 두 방식을 비교한다.

flowchart LR
  H[HTTP 클라이언트] --> R[REST Proxy]
  R -->|produce| T[Kafka 토픽]
  T -->|poll| R
  R --> H

앞의 전송 예시에서는 Content-Type: application/vnd.kafka.json.v2+json이 레코드 값이 JSON이라는 v2 형식을 지정한다. 요청 본문의 records 배열에 레코드가 들어간다. 응답을 받을 때 HTTP 상태만 보지 말고 각 레코드의 partition·offset 또는 오류 정보를 확인한다. 토픽이 존재하지 않거나 브로커가 쓰기를 거부하면 기대한 저장 결과가 나오지 않는다.

v2 소비자는 POST /consumers/{group}으로 인스턴스를 만든 뒤 해당 인스턴스에 토픽을 구독시키고 GET .../records로 읽는다. 생성 응답의 base_uri는 이후 구독·조회·삭제 요청의 기준 경로다. 원문의 source-topic은 이 생성 URL에서 컨슈머 그룹 이름으로 사용됐을 뿐 토픽 이름이라는 뜻은 아니다. 실제 읽을 토픽은 /subscription에 보낸 topics 배열에서 정한다. 컨슈머 인스턴스는 REST Proxy가 유지하는 상태이므로 마지막에는 DELETE로 정리한다. 프록시가 다시 시작되거나 세션이 만료되면 같은 인스턴스 경로가 유효하지 않을 수 있다.

HTTP 단계의미검증
레코드 POST토픽에 값 기록응답의 레코드별 오류·offset 확인
컨슈머 생성그룹의 소비 인스턴스 확보base_uri 저장
subscription읽을 토픽 지정요청 오류와 토픽 이름 확인
records GETKafka poll에 해당빈 배열이면 생산·오프셋 위치 확인
DELETE 인스턴스서버 측 인스턴스 정리이후 같은 경로가 사라지는지 확인

처음부터 메시지가 보이지 않는다면 먼저 프로듀서의 응답, 토픽 이름, auto.offset.reset과 기존 그룹의 커밋된 오프셋을 확인한다. earliest는 그 그룹에 유효한 시작 오프셋이 없을 때 적용되는 설정이므로 이미 읽은 그룹을 무조건 처음으로 되돌리지 않는다. localhost:8082는 REST Proxy를 실행하는 호스트에서 접근할 때만 맞는다. 다른 컨테이너 안에서는 localhost가 그 컨테이너 자신이고, Compose 서비스 이름 등으로 주소를 바꿔야 한다.

v3로 옮길 때의 차이

Confluent의 v3 API는 /v3/clusters/{cluster_id}/topics/{topic_name}/records 같은 경로와 application/json 미디어 타입을 사용한다. v2의 /topics/... 요청을 경로만 바꿔 재사용할 수 없다. 현재 공식 API 참조에서 배포 제품의 독립 REST Proxy와 내장 API 중 어느 것을 쓰는지 확인한다. 특히 v3 Produce API는 HTTP 200이어도 개별 레코드가 실패할 수 있다. 응답의 각 레코드 error_code를 읽어야 조용한 데이터 손실을 피할 수 있다.

REST Proxy 포트를 인터넷에 그대로 노출하면 누가 어떤 토픽에 읽고 쓸 수 있는지 통제하기 어렵다. 실제 서비스에서는 프록시 인증·TLS, Kafka 측 ACL, 요청 크기 제한을 맞추고 클라이언트가 재시도할 때 중복 레코드가 생길 수 있음을 고려한다.

참고: Confluent REST Proxy 개요, API 참조, Quick Start.

같은 카테고리의 글