์นดํ์นด Consumer๋ฅผ ์ฌ์ฉํ๋ค ๋ณด๋ฉด offset์ resetํด์ผํ๋ ๊ฒฝ์ฐ๊ฐ ์ข ์ข ์๋ค.
์ด๋ฐ ๊ฒฝ์ฐ์ Consumer API๋ฅผ ์ฌ์ฉ์ค์ด๋ผ๋ฉด ์ง์ ์ฝ๋ ๋ ๋ฒจ์์ programmaticํ๊ฒ reset๋ ๊ฐ๋ฅํ๊ณ , ์๋๋ฉด kafka์์ ์ ๊ณตํ๋ reset tool์ ์ด์ฉํด์ reset์ด ๊ฐ๋ฅํ๋ค.
์นดํ์นด ๋ฐ์ด๋๋ฆฌ๋ฅผ ์ค์นํ๋ค. Mac์์ brew๊ฐ ์๋ค๋ฉด ์๋ ๋ช
๋ น์ด๋ก ๊ฐ๋จํ ์ค์น ํ ์ ์๋ค.
(Mac์ ๊ธฐ์ค์ผ๋ก ์ค๋ช
ํ๋ค.)
brew install kafka
Kafka 0.10.x ๋ฒ์ ์ดํ์ ๊ฒฝ์ฐ์๋ ์๋ ์ฌ์ฉํ ๋ช ๋ น์ด๋ค์ด ์กด์ฌํ์ง ์๋๋ค. ๋ฐ๋ผ์ ์ค์น๋ kafka ๋ฒ์ ์ด 0.11.x ์ด์์์ ํ์ธํ์. (brew์ ๊ฒฝ์ฐ ํ์ฌ ๊ธฐ๋ณธ์ ์ผ๋ก 2.0 ์ด์์ด ์ค์น๋๋ค.)
consumer group์ ์ง์ ํ๊ณ --describe ์ต์
์ ์ฌ์ฉํ๋ฉด ํ์ฌ consumer group์ offset ์ ๋ณด๋ฅผ ๋ณผ ์ ์๋ค. ๋ช
๋ น์ด๋ ๋ค์๊ณผ ๊ฐ๋ค.
kafka-consumer-groups --bootstrap-server <host:port> --group <group.id> --describe
์คํ๊ฒฐ๊ณผ๋ ๋ค์๊ณผ ๊ฐ๋ค.
TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
example.topic 0 6392623366 6392623859 493 consumer-1-f6f6ffb0-1054-46b9-af13-0b254bc14da0 /10.64.69.95 consumer-1
example.topic 1 6394637143 6394637383 240 consumer-10-6c57b320-7742-4418-8e15-b7d735da346e /10.64.69.95 consumer-2
example.topic 2 6397170269 6397170495 226 consumer-19-dbed41a1-42bb-4ecb-bc8f-84e47c74dbe8 /10.64.69.95 consumer-3
example.topic 3 6397170269 6397170495 226 consumer-19-dbed41a1-42bb-4ecb-bc8f-84e47c74dbe8 /10.64.69.95 consumer-4
๊ฒฐ๊ณผ๋ก ์ถ๋ ฅ๋ ๊ฐ ์ปฌ๋ผ์ ๊ฐ๋จํ๊ฒ ์ค๋ช ํ๋ฉด,
--describe๋ฅผ ํตํด ์กฐํ๋ฅผ ํ์๋ LAG์ด ๊ณ์ ์ผ์ ์์ค์ ์ ์งํ๋ค๋ฉด consumer๊ฐ producer ๊ฐ ๋ง๋ค์ด๋ด๋ ์ด๋ฒคํธ ๋ ์ฝ๋์ ์์ ์ ๋ฐ๋ผ๊ฐ๊ณ ์๋ค๋ ๊ฒ์ ํ์ธํ ์ ์๋ค. ํ์ง๋ง LAG์ด ๊ณ์ ์ฆ๊ฐํ๋ค๋ฉด consumer์ ์ฒ๋ฆฌ ์๋๊ฐ ๋๋ฆฐ ๊ฒ์ด๊ธฐ ๋๋ฌธ์ consumer์ ๊ฐฏ์๋ฅผ ์ถฉ๋ถํ ์ฆ๊ฐ์ํค๊ฑฐ๋, consumer์ ๋ก์ง์ ๋ ๊ฐ๋ตํ ํด์ ๋น ๋ฅธ ์๋๋ก ๋ฐ์ดํฐ ์ฒ๋ฆฌ๋ฅผ ํ ์ ์๋๋ก ๋ณ๊ฒฝํด์ผ ํ๋ค.
kafka์์ ๋ฐ์ดํฐ๋ฅผ ๋ถ๋ฌ์์ ์ฒ๋ฆฌํ๋ ๊ณผ์ ์์ ์ค๋ฅ๊ฐ ๋ฐ์ํ๊ฑฐ๋ ๋ฌธ์ ๊ฐ ๋ฐ๊ฒฌ๋ ๊ฒฝ์ฐ, ๋ค์ ์ํ๋ offset๋ถํฐ ๋ฐ์ดํฐ๋ฅผ ์ฌ์ฒ๋ฆฌ๋ฅผ ํด์ผํ ๊ฒฝ์ฐ๊ฐ ์ข ์ข ์๋ค. ์ด๋ consumer group์ offset reset ๊ธฐ๋ฅ์ ํ์ฉํ๋ฉด ๋๋ค.
์ฃผ์์ฌํญ: consumer group์ด ์คํ์ค์ธ ์ํ์ offset reset์ ์งํํ๋ ๊ฒฝ์ฐ reset์ ์คํจํ๋ค.
kafka-consumer-groups --bootstrap-server <host:port> --group <group> --topic <topic> --reset-offsets --to-earliest --execute
--topic ๋์ --all-topics๋ฅผ ์ง์ ํ๋ฉด ๋ชจ๋ ํ ํฝ์ ๋ํด์ ์คํ์ด ๊ฐ๋ฅํ๋ค.
--execute ์ต์
์ ์ ๊ฑฐํ๊ณ ์คํํ๋ฉด ์ค์ ๋ฐ์๋์ง ์๊ณ ์ด๋ป๊ฒ ๋ณํ ์ง ๊ฒฐ๊ณผ๋ง ์ถ๋ ฅํ๋ dry run์ด ๊ฐ๋ฅํ๋ค.
์คํ์ ์ ์์น๋ฅผ ์ฌ์ค์ ํ๊ธฐ ์ํ ์๋์๊ฐ์ ์์ธ ์ต์ ๋ค์ด ์๋ค.
--shift-by <Long: number-of-offsets> ํ์ (+/- ๋ชจ๋ ๊ฐ๋ฅ)--to-offset <Long: offset>--to-current--by-duration <String: duration> : ํ์ 'PnDTnHnMnS'--to-datetime <String: datetime> : ํ์ 'YYYY-MM-DDTHH:mm:SS.sss'--to-latest--to-earliest--to-datetime์ ๊ฒฝ์ฐ kafka์์ ๋ฐ์ดํฐ๋ฅผ ์ฝ์ด์ ๋ค๋ฅธ๊ณณ์ ์ ์ฅํ๋ ์ค์ ๋ฐ์ดํฐ ์ ์ค ๋๋ ์ค๋ณต write ๋ฑ์ด ๋ฐ์ํ ๊ฒฝ์ฐ์ ๋ ์ง ๋จ์๋ก ๋ฐ์ดํฐ๋ฅผ ๋ค์ ๋ถ๋ฌ์์ ์ฌ์ฒ๋ฆฌํ๊ณ ์ถ์ ๊ฒฝ์ฐ ๋งค์ฐ ์ ์ฉํ๋ค.
Kafka์ ๊ฒฝ์ฐ ์ฌ์ฉ ํํ์ ๋ฐ๋ผ Consumer API์ Stream API ๋๊ฐ์ง๋ฅผ ์ ๊ณตํ๋ค.
Consumer API๋ฅผ ์ฌ์ฉํ ๋ Java์ฝ๋ ๋ ๋ฒจ์์ programmaticalํ๊ฒ offset์ ๋ฆฌ์
ํ๋ ๋ฐฉ๋ฒ์ ๋ค์๊ณผ ๊ฐ๋ค.
๋จผ์ KafkaConsumer๊ฐ ์์ฑํ ํ์
KafkaConsumer<Object, Object> consumer = new KafkaConsumer<>(properties, keyDeser, valueDeser);
consumer loop์ ์ง์ ํ์ฌ consumer.poll()์ ๋ถ๋ฅด๊ธฐ ์ ์, ์์ฑ๋ consumer ๊ฐ์ฒด์ ๋ํด offset์ ๋ณ๊ฒฝํ๋ ๋ค์ ํจ์๋ค์ ํธ์ถํ์ฌ offset์ ์ํ๋ ๋๋ก ์ค์ ํ ์ ์๋ค.
seekToBeginning: earliest๋ก resetseekToEnd: latest๋ก resetseek : ์ง์ offset์ผ๋ก resetoffsetsForTimes: datetime ๊ธฐ์ค์ผ๋ก resethttps://gist.github.com/marwei/cd40657c481f94ebe273ecc16601674b
http://blog.sysco.no/integration/kafka-rewind-consumers-offset/