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

[Kafka] 7. Apache Flume

목차

Apache Flume

원문의 zookeeperConnect, brokerList, topic 설정은 과거 Flume 문서의 방식이다. 현재 Flume User Guide는 Kafka Source/Sink에 kafka.bootstrap.servers, kafka.topics, kafka.topic을 사용한다. 아래 오래된 설정은 기록으로 두고, 새 설정은 글 끝의 예시를 사용한다.

Apache Flume은 오픈소스 프로젝트로 개발된 로그 데이터를 수집 기술이다. 여러 서버에서 생산된 대용량 로그 데이터를 효과적으로 수집하여, HDFS과 같은 원격 목적지에 데이터를 전송하는 기능을 제공한다. 구조가 단순하고 유연하여 다양한 유형의 스트리밍 데이터 플로우(Streaming Data Flow) 아키텍처를 구성할 수 있다.

  • flume은 Source, Channel, Sink 세 가지 요소로 구성
    • 데이터가 추출되는 Source
    • Flume에 데이터를 저장하는 Sink
    • Source에서 Sink or 저장소로 데이터를 전달하는 Channel

Flume 설치

구현

  • source-topic과 target-topic을 생성
  1. flume의 confi 폴더 안에 flume.conf 파일을 생성

flume1을 flume 인스턴스로 선언하고 kafka-source-1, mem-channel-1, kafka-sink-1로 설정한다.

flume1.sources = kafka-source-1
flume1.channels = mem-channel-1
flume1.sinks = kafka-sink-1
## type: 소스타입 설정
flume1.sources.kafka-source-1.type=org.apache.flume.source.kafka.KafkaSource
## zookeeperConnect : 주키퍼의 연결 문자열이며, host:port 형슥올 구분해서 지정한다.
flume1.sources.kafka-source-1.zookeeperConnect = localhost:2181
## topoic 읽어올 소스를 지정한다. flume은 기록할 때 소스당 오직 한 개의 kafka topic을 지원한다.
flume1.sources.kafka-source-1.topic = source-topic
## batchSize : kafka에서 메시지를 가져와서 채널에 기록할 최대 메시지 수이다.(default : 1000) 이 값은 한 번 가져올 때 채널이 처리할 수 있는 데이터양에 따라 결정된다.
flume1.sources.kafka-source-1.batchSize = 100
## channels : source를 연결할 channel 설정
flume1.sources.kafka-source-1.channels = mem-channel-1

## 그외
## batchDurationMillis : 시스템이 배치(batch)를 채널에 기록하기 전에 시스템이 대기할 최대 시간을 밀리초 단위로 지정한다. batchSize가 이시간이 되기전에 초과되면 배치는 해당채널로 전송한다. (default: 1000)
## 메모리가 데이터를 보관하는데 사용되므로 메모리 채널로 설정했다.
## type : 메모리 채널 사용을 설정하기 위해 memory로 설정.
## 채널유형 : memory, JDBC file, kafka channel
flume1.channels.mem-channel-1.type = memory

##그외
## capacity : 메모리에 보관이 가능한 최대 메시지 수. 메모리 용량과 메시지 크기를 고려하여 설정 (default : 100)
## transactionCapacity : 소스에서 가져오거나 저장소에 하나의 트랜잭션으로 처리할 최대 메시지 수

소스의 설정을 선언하는 내용.

source와 저장소 사이의 channel을 정의.

저장소 설정을 선언하는 내용

## type : 저장소유형 지정
flume1.sinks.kafka-sink-1.type = org.apache.flume.sink.kafka.KafkaSink
## brokerList : 메시지를 기록할 kafka cluster broker list이다. host:port 형식으로 여러개 입력가능
flume1.sinks.kafka-sink-1.brokerList = localhost:9092
## 메시지를 기록할 kafka topic
flume1.sinks.kafka-sink-1.topic = target-topic
## 한번에 기록할 메시지 수 지정
flume1.sinks.kafka-sink-1.batchSize = 50
## 데이터를 수집하기 위해 사용할 채널의 이름
flume1.sinks.kafka-sink-1.channel = mem-channel-1

[flume.conf]

flume1.sources = kafka-source-1
flume1.channels = mem-channel-1
flume1.sinks = kafka-sink-1


flume1.sources.kafka-source-1.type=org.apache.flume.source.kafka.KafkaSource
flume1.sources.kafka-source-1.zookeeperConnect = localhost:2181
flume1.sources.kafka-source-1.topic = source-topic
flume1.sources.kafka-source-1.batchSize = 100
flume1.sources.kafka-source-1.channels = mem-channel-1

flume1.channels.mem-channel-1.type = memory

flume1.sinks.kafka-sink-1.type = org.apache.flume.sink.kafka.KafkaSink
flume1.sinks.kafka-sink-1.brokerList = localhost:9092
flume1.sinks.kafka-sink-1.topic = target-topic
flume1.sinks.kafka-sink-1.batchSize = 50
flume1.sinks.kafka-sink-1.channel = mem-channel-1
  1. target-topic에 데이터를 푸시하는 flume 에이전드 실행
$ flume-ng agent --conf-file flume.conf --name flume1

Source → Channel → Sink를 읽는 법

Flume Source는 외부 시스템에서 이벤트를 받아 Channel에 넣고, Sink는 Channel에서 이벤트를 꺼내 목적지로 보낸다. 이 글의 Kafka Source와 Kafka Sink를 한 에이전트에 연결하면 source-topic의 레코드를 target-topic으로 전달한다. 두 토픽이 같으면 자기 출력물을 다시 읽는 루프가 생길 수 있으므로 이름을 분리한다.

flowchart LR
  S[Kafka source-topic] --> F[Flume Kafka Source]
  F --> C[Memory Channel]
  C --> K[Flume Kafka Sink]
  K --> T[Kafka target-topic]

아래는 현재 User Guide의 속성 이름을 사용한 로컬 단일 브로커 실습 설정이다. 실행하는 Flume 릴리스가 해당 Kafka 브로커와 호환되는지는 먼저 확인한다.

flume1.sources = kafka-source-1
flume1.channels = mem-channel-1
flume1.sinks = kafka-sink-1

flume1.sources.kafka-source-1.type = org.apache.flume.source.kafka.KafkaSource
flume1.sources.kafka-source-1.kafka.bootstrap.servers = localhost:9092
flume1.sources.kafka-source-1.kafka.topics = source-topic
flume1.sources.kafka-source-1.kafka.consumer.group.id = flume-forwarder
flume1.sources.kafka-source-1.channels = mem-channel-1

flume1.channels.mem-channel-1.type = memory
flume1.channels.mem-channel-1.capacity = 1000
flume1.channels.mem-channel-1.transactionCapacity = 100

flume1.sinks.kafka-sink-1.type = org.apache.flume.sink.kafka.KafkaSink
flume1.sinks.kafka-sink-1.kafka.bootstrap.servers = localhost:9092
flume1.sinks.kafka-sink-1.kafka.topic = target-topic
flume1.sinks.kafka-sink-1.channel = mem-channel-1

kafka.consumer.group.id는 Source가 자신의 읽기 위치를 관리하는 그룹이다. capacity는 채널이 담을 수 있는 이벤트 수, transactionCapacity는 한 트랜잭션에서 움직이는 최대 수다. Sink가 느리면 채널이 차면서 Source가 더 받기 어려워진다. Memory Channel은 빠르지만 에이전트 프로세스가 종료되면 아직 목적지에 보내지 못한 이벤트가 사라질 수 있다. 내구성이 중요하면 File Channel 등 다른 채널과 디스크·복구 정책을 검토한다.

실행 결과를 어디서 확인할까

Kafka에서 두 토픽을 먼저 만든 뒤 아래처럼 실행하고, 다른 터미널에서 source-topic에 한 줄을 생산한다. Flume이 실행된 상태에서 target-topic 소비자로 같은 줄이 도착하는지 확인한다.

flume-ng agent --conf-file flume.conf --name flume1
# 별도 터미널에서
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic source-topic
# 또 다른 터미널에서
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic target-topic --from-beginning

Flume 로그에는 Source의 Kafka 연결, Channel 용량, Sink 전송 오류가 나타난다. 입력 토픽에는 레코드가 있는데 출력 토픽이 비었다면 Source의 그룹 오프셋, 브로커 주소, 채널 연결 속성, Sink 오류 순으로 좁힌다. localhost는 Flume과 Kafka가 같은 호스트에 있을 때만 맞고 컨테이너 안에서는 각 컨테이너 자신을 가리킨다. Kafka Source는 적어도 한 번 전달 전략이므로 재시작·실패 뒤 중복이 생길 수 있다. 다운스트림에서 이벤트 ID로 중복을 처리해야 할 수 있다.

Flume은 로그·이벤트를 옮기는 도구이지 복잡한 스트림 분석 엔진은 아니다. 단순 Kafka→Kafka 복사만 필요하다면 Flume 프로세스와 채널을 추가할 가치가 있는지 먼저 검토한다. 변환·필터·다른 입력원 수집이 있는 경우 구성 요소를 분리해 관리할 장점이 커진다.

참고: Apache Flume User Guide, Apache Kafka Quickstart.

같은 카테고리의 글