23.06.15
DEVOTEE를 활성화 시키면
지금 작성한 커뮤니티 글에 대해 1개의 댓글을 달아줍니다.
버튼을 누르면 글 수정 시 ChatGPT가 작성한 댓글이 수정됩니다.
| 컨텐츠 유형 | 제목 | 저장일 | 삭제 |
|---|
본인인증 로그인에 실패하였습니다.
회원이 아니시거나 본인인증 등록이
완료되지 않은 사용자입니다.
Kafka는 고성능 데이터 파이프라인, 스트리밍 분석, 데이터 통합 및 미션 크리티컬 어플리케이션을 위한 고성능 분산 이벤트 스트리밍 플랫폼입니다.
Pub-Sub 모델의 메시지 큐 형태로 동작하며, 분산환경에 특화되어 있다는 장점이 존재합니다.
<Fig 1. Publisher-Subscriber Model>
Kafka에 대한 자세한 내용은 아래 문서를 참고하시기 바랍니다.
https://kafka.apache.org/documentation/#gettingStarted
Consumer는 Publisher-Subscriber Model 에서 Subscribe 역할을 담당하는 클라이언트로,
Consumer는 Broker에 저장된 메시지를 가져오는 역할을 수행합니다.
Consumer Group은 개별 Consumer Client을 묶는 논리적인 그룹 단위입니다.
모든 Consumer는 하나의 Consumer Group에 속해야 합니다.
Kafka에서는 특정 Partition에 대한 Offset 관리를 Consumer Group 단위로 수행합니다.
즉, Consumer Group이 여러 개라면, 서로 독립적인 Offset를 갖습니다.
<Fig 2. Consumer Group에 대한 Offset 관리>
즉, 위의 그림과 같이 Consumer Group이 여러 개인 경우,
Kafka의 내부 토픽인 __consumer_offsets를 통해 어떤 Consumer Group이 어떤 토픽의 어떤 파티션의 다음에 읽을 Offset이 몇인지를 기록하여 관리하므로,
서로 다른 Consumer Group 간 Offset은 서로 영향을 줄 수 없습니다.
리밸런싱이란, 특정 토픽을 구독하던 Consumer Group에 변동 사항이 발생했을 때 해당 Group 안에서 파티션을 재분배하는 행위를 의미합니다.
리밸런싱은 아래와 같은 상황에서 발생할 수 있습니다.
Consumer Group 내에 컨슈머가 생성 혹은 삭제된 경우
max.poll.interval.ms로 설정된 시간 내에 poll() 요청을 보내지 못한 경우
session.timeout.ms로 설정된 시간 내에 하트비트를 보내지 못한 경우
Kafka의 토픽은 여러 개의 파티션으로 구성될 수 있습니다.
하나의 파티션은 특정 Consumer Group 안에서 반드시 1개의 Consumer에게 할당되어야 합니다.
<Fig 3. Consumer Group 내 단일 Consumer의 예>
리밸런싱을 간단하게 설명하기 위해, 위의 Fig 3.의 상황을 가장 초기 단계라고 가정하겠습니다.
현재, Consumer Group A에는 Consumer가 1개이므로, 토픽 A의 모든 파티션은 Consumer #1에게 분배됩니다.
<Fig 4. 새로운 컨슈머가 추가되어, 리밸런싱이 발생한 Consumer Group A>
Fig 4.는 기존 Consumer Group A에 새로운 Consumer #2가 추가된 경우입니다.
Consumer #2가 새롭게 추가되어, 해당 Consumer Group 정보로 그룹 코디네이터에 구독 요청을 보내면, 그룹 코디네이터는 해당 Consumer Group에게 리밸런싱할 것을 지시합니다.
이때, 파티션 분배 모드가 EAGER라면 해당 Consumer Group 내부의 모든 컨슈머는 기존 파티션 구독 정보를 모두 버리고 새롭게 파티션을 분배 받게 됩니다.
<Fig 5. 파티션의 개수와 컨슈머의 개수가 동치인 경우>
마지막으로, Fig 5.는 파티션의 개수와 컨슈머의 개수가 일치하는 경우입니다.
모든 파티션이 컨슈머와 1:1 매칭이 되는 구조이기 때문에, 컨슈머의 병렬성을 가장 크게 살릴 수 있는 구조입니다.
마지막으로, Fig 5. 상황에서 기존 컨슈머를 중지하는 경우, 리밸런싱이 어떤 식으로 이루어지는지 알아보도록 하겠습니다.
(venv6) [ec2-user@kcluster22 chapter6]$ /usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server kcluster23.foo.bar:9092 --group k-consumer01 --describe
OpenJDK 64-Bit Server VM warning: If the number of processors is expected to increase from one, then you should configure the number of parallel GC threads appropriately using -XX:ParallelGCThreads=N
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
k-consumer01 chapter06 1 132 132 0 rdkafka-18d475ee-1556-4606-b303-5a9f5ebfee81 /192.168.100.53 rdkafka
k-consumer01 chapter06 2 96 96 0 rdkafka-52b13372-5521-4e6f-b997-9d8616fe9088 /192.168.100.32 rdkafka
k-consumer01 chapter06 0 72 72 0 rdkafka-efb00444-ffdb-4b92-9165-a5f07755b2cd /192.168.100.58 rdkafka위의 콘솔 화면은 모든 파티션과 Consumer Group 내의 컨슈머가 1:1으로 매칭된 상태입니다.
이 상황에서 가장 첫 번째 컨슈머인 rdkafka-18d475ee-1556-4606-b303-5a9f5ebfee81를 종료하겠습니다.
(venv6) [ec2-user@kcluster22 chapter6]$ /usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server kcluster23.foo.bar:9092 --group k-consumer01 --describe
OpenJDK 64-Bit Server VM warning: If the number of processors is expected to increase from one, then you should configure the number of parallel GC threads appropriately using -XX:ParallelGCThreads=N
Warning: Consumer group 'k-consumer01' is rebalancing.
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
k-consumer01 chapter06 0 132 132 0 - - -
k-consumer01 chapter06 2 72 72 0 - - -
k-consumer01 chapter06 1 96 96 0 - - -Warning 메시지에서 보이듯, k-consumer01이라는 Consumer Group에 리밸런싱이 발생했다는 메시지가 출력되며, 기존 파티션 할당 정보가 모두 삭제된 것을 확인할 수 있습니다.
(venv6) [ec2-user@kcluster22 chapter6]$ /usr/local/kafka/bin/kafka-consumer-groups.sh --bootstrap-server kcluster23.foo.bar:9092 --group k-consumer01 --describe
OpenJDK 64-Bit Server VM warning: If the number of processors is expected to increase from one, then you should configure the number of parallel GC threads appropriately using -XX:ParallelGCThreads=N
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
k-consumer01 chapter06 0 132 132 0 rdkafka-52b13372-5521-4e6f-b997-9d8616fe9088 /192.168.100.58 rdkafka
k-consumer01 chapter06 1 96 96 0 rdkafka-52b13372-5521-4e6f-b997-9d8616fe9088 /192.168.100.58 rdkafka
k-consumer01 chapter06 2 72 72 0 rdkafka-efb00444-ffdb-4b92-9165-a5f07755b2cd /192.168.100.32 rdkafka잠시 후, 원래 두 번째 파티션만 할당 받았던 rdkafka-52b13372-5521-4e6f-b997-9d8616fe9088가 첫 번째 파티션도 추가적으로 할당 받은 것을 확인할 수 있습니다.
Note
Fig 5.에서 언급한 바와 같이, 파티션과 컨슈머의 개수가 일치할 때 병렬성이 가장 좋다고 했습니다.
그렇다면, 파티션보다 컨슈머의 개수가 더 많을 수록 병렬성이 더욱 증가하지 않을까요?
정답은 No 입니다. 파티션보다 컨슈머의 개수가 더 많다면, 초과한 컨슈머에 대해서는 파티션을 할당 받지 못해, 해당 컨슈머는 Stand-by 상태로 남아 있게 됩니다.
오히려, 파티션을 할당 받지 못한 컨슈머는 Broker의 TCP Connection을 불필요하게 낭비합니다.
Warning
Kafka에서 리밸런싱이 발생할 때, 해당 Consumer Group은 리밸런싱이 완료될 때까지 파티션의 메시지를 읽을 수 없습니다.
즉, 리밸런싱은 되도록이면 발생하지 않도록 하는 것이 좋으며 이 현상을 Stop the world라고 표현합니다.
<Fig 6. 컨슈머의 subscribe(), poll(), commit()>
subscribe()
컨슈머가 Consumer Group의 정보를 이용하여, 브로커에게 구독 요청을 하는 메서드 입니다.
브로커는 해당 요청을 받으면, 해당 Consumer Group에 대해 리밸런싱 지시를 내도록 합니다.
poll()
주기적으로, 할당 받은 파티션에서 메시지를 가져오는 역할을 수행하는 메서드입니다.
특이점으로는, 첫 번째 poll()에서는 메시지를 가져오지 않고 브로커의 metadata를 가져오고, 그룹 코디네이터와 연결하는 작업을 수행합니다.
즉, 두 번째 poll()에서부터 메시지를 가져옵니다.
commit()
브로커의 내부 토픽인 __consumer_offsets에 특정 토픽의 특정 파티션의 어떤 offset 까지 읽었는지 기록합니다.
본 메서드는 개발자가 직접 manual commit을 할 수 있으며, 주기적으로 auto commit을 하도록 할 수도 있습니다.
<Fig 7. 컨슈머의 내부 구성 요소>
Fetcher & ConsumerNetworkClient
2개의 컴포넌트를 통해, 파티션의 데이터를 해당 컨슈머 클라이언트로 가져옵니다.
ConsumerCoordinator
해당 Consumer Group의 리더 컨슈머가 누구인지, 해당 컨슈머의 옵션 등을 관리합니다.
HeartBeat Thread
하트비트 체크를 위한 별도의 스레드입니다.
SubscriptionState
파티션의 구독 상태를 관리합니다.
Fetcher와 ConsumerNetworkClient는 파티션의 데이터를 해당 컨슈머로 가져오는 기능을 수행합니다.
<Fig 8. Fetcher & ConsumerNetworkClient>
Fig 8.에서 보이듯, 실제 파티션에서 데이터를 받아오는 기능은 ConsumerNetworkClient가 수행합니다.
public class ConsumerNetworkClient implements Closeable {
private static final int MAX_POLL_TIMEOUT_MS = 5000;
private final KafkaClient client;
private final Metadata metadata;
private final long retryBackoffMs;
private final int maxPollTimeoutMs;
private final int requestTimeoutMs;
private final AtomicBoolean wakeupDisabled = new AtomicBoolean();
private final ReentrantLock lock = new ReentrantLock(true);
private final ConcurrentLinkedQueue<RequestFutureCompletionHandler> pendingCompletion = new ConcurrentLinkedQueue<>();
private final ConcurrentLinkedQueue<Node> pendingDisconnects = new ConcurrentLinkedQueue<>();
...
}ConsumerNetworkClient는 비동기로 동작하여, 파티션에서 데이터를 받아온 후 pendingCompletion에 bytes 형태로 저장합니다.
ConsumerNetworkClient는 비동기로 데이터를 파티션에서 받아 오기도 하지만, Fetcher가 추가적인 Request를 생성하여 데이터 요청을 할 수도 있습니다.
public class Fetcher<K, V> implements Closeable {
private final ConsumerNetworkClient client; // ConsumerNetworkClient를 멤버 변수로 사용합니다.
private final int minBytes;
private final int maxBytes;
private final int maxWaitMs;
private final int fetchSize;
private final long retryBackoffMs;
private final long requestTimeoutMs;
private final int maxPollRecords;
private final ConsumerMetadata metadata;
private final SubscriptionState subscriptions;
private final ConcurrentLinkedQueue<CompletedFetch> completedFetches;
private final BufferSupplier decompressionBufferSupplier = BufferSupplier.create();
private final Deserializer<K> keyDeserializer;
private final Deserializer<V> valueDeserializer;
private final IsolationLevel isolationLevel;
...
}Fetcher는 ConsumerNetworkClient를 멤버 변수로 가지고 있으며,
주기적으로 ConsumerNetworkClient가 받아온 bytes 형태의 데이터를 꺼내,
keyDeserializer 및 valueDeseializer를 통해 역-직렬화를 수행한 후 completedFetches에 삽입합니다.
엄밀히 말하면, max.poll.records는 파티션에서 읽어올 메시지 개수가 아닌, completedFetches (Fig 7.에서 In-memory Linked Queue를 의미.)에서 읽어올 최대 개수를 의미합니다.
HeartBeat Thread는 그룹 코디네이터에게 해당 컨슈머의 정상적인 활동을 보고하는 별도의 스레드입니다.
해당 스레드는 첫 번째 poll() 메서드가 호출될 때 생성됩니다.
<Fig 9. HeartBeat Thread의 예>
Fig 9.처럼 별도의 스레드가 그룹 코디네이터에게 주기적으로 HeartBeat API를 보내는 것을 확인할 수 있습니다.
하트비트 관련한 옵션은 아래와 같습니다.
heartbeat.interval.ms = 하트비트 스레드가 HeartBeat API를 보내는 간격을 의미합니다.
이 값은 통상적으로 session.timeout.ms의 1/3 수준으로 설정될 것을 권장합니다.
session.timeout.ms = 하트비트 타임아웃 값으로, 해당 시간 내에 그룹 코디네이터가 하트비트를 받지 못하면 해당 Consumer Group에 리밸런싱 명령을 내립니다.
max.poll.interval.ms = 이전 poll() 호출 후, 다음 poll() 호출까지 그룹 코디네이터가 기다리는 시간입니다.
이 값 이내에 poll() 요청을 보내지 않는 경우, 해당 Consumer Group에 리밸런싱 명령을 내립니다.
그룹 코디네이터는 브로커 내부에 위치한 컴포넌트로 Consumer Group의 상태를 체크하는 역할을 담당합니다.
특정 Consumer Group의 상태를 체크하여, 그룹 내 변동이 발생하거나 특정 컨슈머에 장애가 발생한 경우 리밸런싱 명령을 내립니다.
아래는 Consumer Group A에 컨슈머 C#3이 새롭게 합류한 상황에서의 그룹 코디네이터와 Consumer Group이 통신하는 흐름을 다룬 그림입니다.
<Fig 10. 그룹 코디네이터 리밸런싱 흐름>
FindCoordinator : C#3가 시작되어, Consumer Group에 변동이 생겨, 리밸런싱을 시작해야 하는 상태입니다.
FindCoordinator Request를 통해 그룹 코디네이터의 정보를 얻어 옵니다.
JoinGroup Request : 모든 컨슈머는 JoinGroup Request를 그룹 코디네이터에게 보내 Rebalance Protocol를 초기화 합니다.
그룹 코디네이터는 어떤 컨슈머를 Group에 참여시킬지 결정합니다.
(session.timeout.ms 시간 내에 HeartBeat가 도착한 Consumer만 참여할 수 있습니다.)
JoinGroup Response : 가장 먼저 JoinGroup Request를 보낸 컨슈머는 Leader Consumer가 되고, 응답 값으로 활성화된 컨슈머 목록과 파티션 분배 전략을 받습니다.
그 외의 컨슈머들은 응답 값으로 빈 값을 받습니다.
SyncGroup Request & Response : 모든 컨슈머가 SyncGroup Request를 보냅니다.
여기서, Leader Consumer는 파티션 분배 전략에 맞게 분배된 파티션 매칭 정보를 추가로 담습니다.
그룹 코디네이터는 모든 컨슈머에게 요청을 받으면, Leader Consumer가 계산한 정보를 각 컨슈머에게 담당 파티션을 담아 응답합니다.
Fetching : SyncGroup Response로부터 받은 파티션 정보를 통해, 메시지를 구독합니다.
특정 Consumer Group에 리밸런싱이 발생하는 경우, 관련 컨슈머들은 해당 시간 동안 모든 구독 활동이 중지되기 때문에 메시지를 구독할 수 없습니다.
위에서 이를 Stop the world 현상이라고 소개했습니다.
Static Membership이란, 컨슈머에게 고정 ID를 부여하여, 특정 컨슈머의 시스템 업데이트, Consumer Parameter 조정 등의 이슈로 재시작 하는 경우 불필요한 리밸런싱을 줄이기 위해 사용됩니다.
<Fig 11. Static Membership - group.instance.id>
Fig 11.에서 보이는 바와 같이, Static Membership은 Consumer Group 내부의 각 컨슈머에 대한 고유 ID를 group.instance.id를 통해 지정함으로써 설정할 수 있습니다.
group.instance.id가 null이 아니라면, 해당 컨슈머는 프로세스를 종료해도 그룹 코디네이터에게 종료 사실을 알리지 않습니다.
<Fig 12. Static Membership이 적용되지 않은 Consumer Group과 적용된 Consumer Group의 예>
Fig 12.의 Consumer Group A는 Static Membership이 적용되지 않았을 때의 상황입니다.
Consumer #3이 종료가 될 때, 종료 사실을 그룹 코디네이터에게 알리고, 그룹 코디네이터는 Leader Consumer에게 리밸런싱 명령을 내립니다.
그와 반대로, Consumer Group B는 group.instance.id에 의해 Static Membership이 적용되어 있을 때입니다.
이 경우, consumer_group_b-consumer-3이 종료될 때, 종료 사실을 그룹 코디네이터에게 알리지 않으므로, 리밸런싱이 일어나지 않습니다.
consumer_group_b-consumer-3는 session.timeout.ms 시간 내에 다시 기동되어, Heart beat를 보낸다면, 리밸런싱 없이 다시 시작될 수 있습니다.
리밸런싱이 발생할 때, Leader Consumer는 그룹 코디네이터에게 파티션 분배 전략을 받고, 각 컨슈머와 파티션을 매칭합니다.
컨슈머의 파티션 할당 전략은 아래와 같이 크게 4가지 방법이 존재합니다.
파티션 할당 전략 | 설명 | Rebalancing Protocol |
|---|---|---|
Range Partition 할당 전략 | 기본 값으로, Topic 별로 동일한 Partition을 특정 Consumer에게 할당하는 방식입니다. | EAGER |
Round Robin Partition 할당 전략 | 사용 가능한 Partition과 Consumer를 순차적으로 할당합니다. | EAGER |
Sticky Partition 할당 전략 | Consumer가 구독 중인 Partition을 계속 유지하게 끔 할당합니다. | EAGER |
Cooperative Sticky Partition 할당 전략 | Sticky와 유사하지만, 전체 Rebalancing이 아닌 필요한 Partition끼리 점진적으로 Rebalancing 하는 방식입니다. | COOPERATIVE |
<Fig 13.EAGER 모드의 파티션 할당 전략의 예>
Fig 13.은 Consumer Group A에 Consumer #3이 추가될 때의 상황에서 파티션 할당 전략이 EAGER 모드일 때의 순서입니다.
EAGER 모드는 Consumer Group 내의 모든 컨슈머가 리밸런싱이 발생합니다.
<Fig 14. COOPERATIVE 모드의 파티션 할당 전략의 예>
Fig 14.는 똑같은 상황에서 Cooperative 모드의 예입니다.
이 경우에는 대상이 되는 Consumer에 대해서만 Rebalancing이 발생합니다.
서로 다른 토픽들의 동일한 파티션을 같은 컨슈머에게 할당하는 방식입니다.
다만, 균등하게 Partition이 분배되지 않으므로, 분배가 불균형할 수 있다는 것을 유의해야 합니다.
이 전략은 특히, Key가 존재하는 메시지를 다루는 경우에 유용합니다.
예를 들어, 주문 처리 시스템에서 토픽 A는 order_id를 관리하고 토픽 B에서는 order_id에 대한 order_item를 관리하는 경우,
서로 다른 토픽이라고 하더라도 하나의 컨슈머에 동일한 파티션이 분배되기 때문에 데이터 처리가 용이할 수 있습니다.
(Round Robin 방식이라면, order_id를 받아 DB에 조회하는 로직이 필수적이지만, 해당 전략의 경우에는 그 빈도를 낮출 수 있습니다.)
<Fig 15. Range Partition 할당 전략의 예.>
해당 전략은 파티션을 순차적으로 컨슈머에게 할당되어 파티션 분배가 비교적 균등하게 할당될 수 있습니다.
Round Robin을 위해, 모든 파티션을 순서대로 배치 후에 차례대로 할당합니다.
<Fig 16. Round Robin Partition 할당 전략의 예>
Range Partion 할당 전략과 Round Robin Partition 할당 전략은 리밸런싱이 일어나면 동일한 컨슈머와 파티션이 분배되는 것을 보장할 수 없습니다.
Sticky Partition 할당 전략은 두 가지 목적으로 할당 전략을 제공합니다.
최대한 균형 있게 파티션을 분배한다.
리밸런싱이 발생할 때 되도록 기존 할당된 파티션이 할당되도록 노력한다.
위에서 “노력”이라는 단어에서 볼 수 있듯이, Sticky Partition 할당 전략을 사용해도 기존에 할당 받은 파티션이 할당되는 것을 무조건 보장할 수는 없습니다.
Sticky Partition 할당 전략의 초기 동작은 Round Robin과 매우 흡사합니다.
<Fig 17. Sticky Partition 할당 전략의 예>
Fig 17.의 왼쪽 상황은 가장 초기에 Sticky Partition의 할당 결과입니다. Round Robin과 매우 유사한 형태로 균등하게 파티션 분배가 완료된 것을 확인할 수 있습니다.
Fig 17.의 오른쪽 상황은 Consumer #2에 장애가 생겨, Group에서 제외되고 리밸런싱이 된 후의 상황입니다.
파란색과 빨간색 선에서 보이듯이, Sticky Partition 할당 전략은 기존 매핑 정보를 이용하여 원래의 매핑 정보를 최대한 지키려고 노력합니다.
즉, 기존의 매핑 정보가 존재하는 경우를 먼저 할당하며, 매핑 정보가 일치하지 않는 케이스에 대해서 균등하게 파티션을 재분배합니다.
기존 Rebalancing Protocol인 EAGER는 Consumer Group 내에 변경이 감지되면, 모든 컨슈마가 리밸런싱이 수행된다는 특징이 있습니다.
EAGER를 쓰는 이유는 2가지 특징이 존재합니다.
Consumer 간의 파티션 소유권 변경 이슈 = Consumer A에서 Consumer B로 Partition #1의 소유권을 넘기기 위해, 모든 소유권을 삭제
로직의 단순화
EAGER의 경우, 리밸런싱이 발생할 때 모든 컨슈머는 메시지를 구독할 수 없는 Stop the world 현상이 발생합니다.
이에 반해, 프로듀서의 경우에는 리밸런싱에 영향 없이 파티션에 데이터를 계속해서 쓰기 때문에 다운 타임 동안에는 LAG가 급격하게 증가하게 됩니다.
이러한 이슈를 개선하기 위해, 전체 컨슈머에 대해서 리밸런싱하는 것이 아닌, 필요한 파티션에 대해서만 리밸런싱을 점진적으로 하는 방식인 COOPERATIVE가 등장했습니다.
<Fig 18.Cooperative Sticky Partition 할당 전략의 예>
Fig 18.은 Cooperative Sticky Partion 할당 전략의 예입니다. 위의 그림은 새로운 컨슈머가 추가될 때이고, 아래는 기존의 컨슈머가 제외될 때의 예입니다.
Fig 17.과 다른 점은 2번째 Step에서 모든 컨슈머의 파티션 구독 정보를 초기화 하지 않고, 대상 파티션만 리밸런싱하는 특징을 가지고 있습니다.
프로듀서의 중복 없이 전송 관련된 개념을 다룰 때, Transaction Coordinator가 등장하며
프로듀서는 해당 컴포넌트와 통신하며 메시지에 대해 PID와 메시지 sequential number를 활용하여, Transaction을 부여하는 방법으로 중복 없이 전송을 지원했습니다.
컨슈머의 입장에서 정확히 한 번 읽기란, Tranaction이 완료된 메시지만 읽는 것을 의미합니다.
<Fig 19. Kafka Transaction Flow>
Fig 19.는 프로듀서가 발행한 메시지에 대한 흐름을 다룹니다.
produce 이후의 파티션 상태를 보면 빨간색 바탕으로 a라는 특수하게 발행된 메시지와 연두색 바탕의 c라는 메시지를 확인할 수 있습니다.
빨간색 바탕의 a는 Tranaction Producer가 발행한 Commit이 성공한 메시지임을 나타내는 특수 표시이고, 그 뒤에 연두색 바탕의 c는 Transaction Commit Message 입니다.
즉, 컨슈머에서 정확히 한 번 읽기 옵션을 사용하는 경우, 빨간색 바탕의 a에 대한 메시지만 읽는 것을 뜻합니다.
(Trasaction Commit Message는 일반적으로 Transaction 종료 표시를 위해 남기는 메시지입니다. 해당 메시지 때문에, Transaction Producer는 메시지를 1개만 발행해도 offset은 2개가 증가합니다.)
컨슈머 클라이언트에서 아래와 같은 설정을 통해, 정확히 한 번 읽기 기능을 사용할 수 있습니다.
enable.auto.commit = false
isolation.level = read_committed
여기서, enable.auto.commit = false로 설정한 이유는 특정 시간마다, 주기적으로 commit 하기 때문에 읽은 메시지가 정상적으로 처리되었는지 판단하기 어렵습니다.
주로, enable.auto.commit = true를 설정하면 At Most Once로 동작하게 됩니다.
컨슈머의 경우, Tranaction Producer와 달리 Tranaction Coordinator와 통신하지 않기 때문에 정확하게 메시지를 한 번 가져오는 것은 보장할 수 없습니다.
Transaction 범위를 컨슈머의 동작까지 Exactly-Once를 구현하기 위해서는 프로듀서에서 send() API를 사용하는 것이 아닌,
sendOffsetsToTransaction() API를 호출하여 Consumer Group가 소비한 Offset Commit을 Transaction 범위에 포함해야 합니다.
kafkaProducer.sendOffsetsToTransaction(consumedOffsets, "kafka-transactions-group")
// consume된 offsets과 대상 Consumer Group ID를 파라미터로 넣습니다.Transaction에 특정 Consumer Group이 소비한 offset를 추가로 넣어 Transaction Coordinator에게 전달하기 때문에,
특정 컨슈머가 가지고 있는 offset과 비교하는 방식으로 Transaction 성공 유무를 확인하기 때문에, Producer-Consumer 간의 정확히 한 번 전송을 지원할 수 있습니다.
본 문서는 KRU의 2023 KAFKA 온라인 스터디의 후원을 통해 작성되었습니다.
[1] https://www.aladin.co.kr/shop/wproduct.aspx?ItemId=281606911
[3] https://jinhanchoi1.medium.com/kafkaconsumer-and-tcp-connection-5a2c8b197732
[4] https://ojt90902.tistory.com/1092
DEVOTEE를 활성화 시키면
지금 작성한 댓글에 AI가 댓글을 달아줍니다.