메뉴 건너뛰기

Bigdata, Semantic IoT, Hadoop, NoSQL

Bigdata, Hadoop ecosystem, Semantic IoT등의 프로젝트를 진행중에 습득한 내용을 정리하는 곳입니다.
필요한 분을 위해서 공개하고 있습니다. 문의사항은 gooper@gooper.com로 메일을 보내주세요.


* 출처 : http://epicdevs.com/21

* 샘플 소스 : https://www.gooper.com/ss/index.php?mid=bigdata&category=2776&document_srl=3161


<Producer>

필수 프로퍼티

 프로퍼티 

설명 

 metadata.broker.list

메타데이터를 받아 올 Kafka broker 리스트. 호스트1:포트1,호스트2:포트2,호스트3:포트3의 형태로 명시한다.
예) kafka-test-001.epicdevs.com:9092,kafka-test-002.epicdevs.com:9092,kafka-test-003.epicdevs.com:9092


여기서 명시하는 broker는 메타데이터를 받아오는 데만 사용하고, 실제 메시지를 전송할 때에는 메타데이터를 기반으로 새로운 connection을 맺은 다음 메시지를 전송한다. 따라서 이 리스트에는 전체 broker 중 일부만 명시해도 무관하다.


중요 프로퍼티

 프로퍼티

 기본 값 

설명 

 serializer.class

 kafka.serializer.DefaultEncoder

메시지를 serialize할 때 사용하는 인코더.DefaultEncoder는 byte[]를 받아서 그대로 전달한다.

 key.serializer.class

 serializer.class의 값과 동일

메시지 키를 serialize할 때 사용하는 인코더.

 partitioner.class

 kafka.producer.DefaultPartitioner

메시지를 어떤 partition에 전송할지 결정하는 클래스. DefaultPartitioner는 메시지 키의 해시 코드를 기반으로 메시지를 전송할 partition을 결정한다. 메시지 키를 명시하지 않았거나 null 값을 키로 전달할 경우 사용자가 명시한 partitioner.class를 무시하고 임의의 partition에 메시지를 보내게 된다. 이 때문에 특정 상황에서 전체 partition에 메시지가 제대로 분산되지 않는 현상이 발생할 수 있다. 이에 대한 자세한 사항은 Kafka FAQ의 Why is data not evenly distributed among partitions when a partitioning key is not specified?를 참조하길 바란다.

 request.required.acks

 0

Producer가 전송한 메시지가 몇 개의 replica에 commit되어야 ack처리(성공적으로 전송된 것으로 간주)를 하는지 결정하는 기준.

  • 0: producer는 broker로부터 ack를 기다리지 않고 메시지 전송이 끝나자마자 성공된 것으로 간주한다. 응답 시간은 가장 빠르지만 broker에서 오류가 발생할 경우 메시지가 유실된다.
  • 1: leader를 맡고있는 replica에 메시지가 commit되면 ack처리를 한다.
  • N: N개의 replica에 메시지가 commit되면 ack처리를 한다.
  • -1: 모든 replica에 메시지가 commit되면 ack처리를 한다.

 compression.codec

 none

메시지를 압축할 때 사용할 코덱. nonegzip,snappy 중 하나를 선택할 수 있다. none을 선택하면 메시지를 압축하지 않는다.

 producer.type

 sync

Producer가 메시지를 동기적으로 전송할지 비동기적으로 전송할지에 대한 설정. 동기적으로 전송하려면 sync로 명시하고 비동기적으로 전송하려면async로 명시한다. 비동기 producer를 사용할 경우 메시지를 일정 시간 큐에 쌓아 두었다가 한 번에 전송하므로 producer의 메시지 처리량을 향상시킬 수 있다. 단, 장애가 발생할 경우 전송하지 않고 쌓아 둔 메시지가 유실될 우려가 있다.

 queue.buffering.max.ms

 5000

비동기 producer를 사용할 경우 몇 ms동안 메시지를 모아둘지 결정하는 값. 비동기 producerqueue.buffering.max.ms 값에 도달하거나batch.num.messages 값에 도달할 경우 쌓아두었던 메시지를 전송한다.

 batch.num.messages

 200

비동기 producer를 사용할 경우 최대 몇 개의 메시지를 모아둘지 결정하는 값. 비동기 producer는queue.buffering.max.ms 값에 도달하거나batch.num.messages 값에 도달할 경우 쌓아두었던 메시지를 전송한다.

위에서 언급한 필수 프로퍼티와 중요 프로퍼티 외의 항목들은 Kafka 공식 페이지의 3.3 Producer Configs를 참고하길 바란다.



<Consumer>

코드 상에는 consumer가 소비할 메시지의 offset과 관련된 내용은 전혀 존재하지 않는다. Offset 값은 Zookeeper에서 별도로 관리하며, high level consumer는 Zookeeper로부터 자신이 속한 consumer group이 몇 번째 메시지 offset을 소비할 차례인지 전달받은 뒤 해당 offset의 메시지부터 소비하기 시작한다.

필수 프로퍼티

 프로퍼티 

설명 

 group.id

Consumer가 속한 consumer group의 ID. Zookeeper에서는 각 consumer group의 메시지 offset을 관리하는데, 이 때 이 ID가 키로써 사용된다. 따라서 consumer group ID가 같으면 모두 같은 consumer group에 속한 것으로 간주되며 메시지 offset 값 또한 공유된다.

 zookeeper.connect

Zookeeper 인스턴스 리스트. 호스트1:포트1,호스트2:포트2,호스트3:포트3의 형태로 명시한다.
예) kafka-test-001.epicdevs.com:2181,kafka-test-002.epicdevs.com:2181,kafka-test-003.epicdevs.com:2181


중요 프로퍼티

 프로퍼티

 기본 값 

설명 

 auto.commit.enable true

Consumer가 메시지를 전달받은 뒤 자동으로 offset 값을 commit할지 결정하는 플래그. 메시지가 성공적으로 처리되었을 때만 offset이 commit되도록 하려면 이 값을 false로 설정해야 한다. 이 값이 true일 경우auto.commit.interval.ms 값마다 주기적으로 offset을 commit하며,false일 경우 ConsumerConnector의 commitOffsets 메소드를 직접 호출해야 offset이 commit된다

 auto.commit.interval.ms 60000

auto.commit.enable이 true일 때 offset을 자동으로 commit하는 주기. 이 값을 길게 잡으면 메시지 처리 중에 장애가 발생할 경우 실제로 처리된 메시지 offset과 commit된 offset 간의 격차가 커져서 fail over 후 중복으로 처리되는 메시지의 수가 많아질 가능성이 높으며, 짧게 잡을 경우 잦은 Zookeeper 업데이트로 인한 오버헤드가 발생할 수 있다.

 auto.offset.reset largest

Consumer가 속한 consumer group의 offset 값이 존재하지 않거나 범위를 벗어나는 값을 전달받았을 경우 어떻게 동작할지를 정하는 값.

  • smallest: 가장 작은 offset의 메시지부터 소비한다.
  • largest: 가장 큰 offset의 메시지 이후부터 소비한다. (즉, 새롭게 전송되는 메시지부터 소비한다.)

위에서 언급한 필수 프로퍼티와 중요 프로퍼티 외의 항목들은 Kafka 공식 페이지의 3.2 Consumer Configs를 참고하길 바란다.

번호 제목 글쓴이 날짜 조회 수
740 bananapi 5대(ubuntu계열 리눅스)에 yarn(hadoop 2.6.0)설치하기-ResourceManager HA/HDFS HA포함, JobHistory포함 총관리자 2015.04.24 19143
739 mapreduce appliction을 실행시 "is running beyond virtual memory limits" 오류 발생시 조치사항 총관리자 2017.05.04 16898
738 org.apache.hadoop.hdfs.server.common.InconsistentFSStateException: Directory /tmp/hadoop-root/dfs/name is in an inconsistent state: storage directory does not exist or is not accessible. 구퍼 2013.03.11 14781
737 drop table로 삭제했으나 tablet server에는 여전히 존재하는 테이블 삭제방법 총관리자 2021.07.09 7554
736 insert hbase by hive ... error occured after 5 hours..HMaster가 뜨지 않는 장애에 대한 복구 방법 총관리자 2014.04.29 7129
735 Resource temporarily unavailable(자원이 일시적으로 사용 불가능함) 오류조치 총관리자 2015.11.19 6854
734 HBase shell로 작업하기 구퍼 2013.03.15 5834
733 dr.who로 공격들어오는 경우 조치방법 file 총관리자 2018.06.09 5603
732 하둡 분산 파일 시스템을 기반으로 색인하고 검색하기 구퍼 2013.03.15 5573
731 [Decommission]시 시간이 많이 걸리면서(수일) Decommission이 완료되지 않는 경우 조치 총관리자 2018.01.03 5320
730 Ubuntu 16.04LTS 설치후 초기에 주어야 하는 작업(php, apache, mariadb설치및 OS보안설정등) file 총관리자 2017.05.23 5270
729 hive 2.0.1 설치및 mariadb로 metastore 설정 총관리자 2016.06.03 5184
728 Hive Query Examples from test code (2 of 2) 총관리자 2014.03.26 5012
727 Spark에서 Serializable관련 오류및 조치사항 총관리자 2017.04.21 4901
726 [gson]mongodb의 api를 이용하여 데이타를 가져올때 "com.google.gson.stream.MalformedJsonException: Unterminated object at line..." 오류발생시 조치사항 총관리자 2017.12.11 4411
725 import 혹은 export할때 hive파일의 default 구분자는 --input-fields-terminated-by "x01"와 같이 지정해야함 총관리자 2014.05.20 4245
724 checking for termcap functions library... configure: error: No curses/termcap library found 구퍼 2013.03.08 4120
723 sqoop작업시 hdfs의 개수보다 더많은 값이 중복되어 oracle에 입력되는 경우가 있음 총관리자 2014.09.02 4093
722 다수의 로그 에이전트로 부터 로그를 받아 각각의 파일로 저장하는 방법(interceptor및 multiplexing) 총관리자 2014.04.04 4089
721 .git폴더를 삭제하고 다시 git에 추가하고 서버에 반영하는 방법 총관리자 2017.06.19 4077

A personal place to organize information learned during the development of such Hadoop, Hive, Hbase, Semantic IoT, etc.
We are open to the required minutes. Please send inquiries to gooper@gooper.com.

위로