Flume Integration with Kafka

Below we will see how kafka can be integrated with Flume as a Source, Channel and Sink. If you are new to Flume and Kafka, you can refer FLUME and KAFKA.

Let us Start :

  1. Using Kafka as a SOURCE for Flume:
    We want to pass messages to a Kafka Producer, which will go through Flume channel (In-Memory) and finally getting Stored in Flume Sink (say HDFS).
    Start Zookeeper & Kafka Services and create a Topic ‘testing’ :

     

    nohup ./${kafka_home}/bin/zookeeper-server-start.sh ${kafka_home}/conf/zookeeper.properties &
    nohup ./${kafka_home}/bin/kafka-server-start.sh ${kafka_home}/conf/server.properties &
    ./${kafka_home}/bin/kafka-topic.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic testing
    ./${kafka_home}/bin/kafka-topic.sh --list --zookeeper localhost:2181

    Set-up Flume Conf:

    vi ${flume_home}/conf/excercise5.conf and add below contents:
    ## Configuring Components
    a1.sources=source1
    a1.channels=channel1
    a1.sinks=sink1
    ## Configuring Source
    a1.sources.source1.type=org.apache.flume.source.kafka.KafkaSource
    a1.sources.source1.zookeeperConnect=localhost:2181
    a1.sources.source1.topic=testing
    a1.sources.source1.groupId=flume
    a1.sources.source1.channels=channel1
    a1.sources.source1.interceptors=i1
    a1.sources.source1.interceptors.i1.type=timestamp
    a1.sources.source1.kafka.consumer.timeout.ms=100
    ## Configuring Channel
    a1.channels.channel1.type=memory
    a1.channels.channel1.capacity=10000
    a1.channels.channel1.transactionCapacity=1000
    ## Configuring sink
    a1.sinks.sink1.type=hdfs
    a1.sinks.sink1.channel=channel1
    a1.sinks.sink1.hdfs.path=/tmp/kafka/%{topic}/%y-%m-%d
    ## Number of seconds to wait before rolling current file
    a1.sinks.sink1.hdfs.rollInterval=5
    ## File size to trigger roll, in byte
    a1.sinks.sink1.hdfs.rollSize=1024
    ## Number of events written to file before it rolled
    a1.sinks.sink1.hdfs.rollCount=10
    ## DataStrem - Stores data instead of ASCII values of data
    a1.sinks.sink1.hdfs.fileType=DataStream

    Start Flume agent:

    ./${flume_home}/bin/flume-ng agent -c ${flume_home}/conf/ -f ${flume_home/conf/exercise5.conf --name a1 -Dflume.root.logger=INFO,console

    Pass your Messages from Kafka Producer and Monitor the Flume agent window:

    ./${kafka_home}/bin/kafka-console-producer.sh --broker-list localhost:9092 --topic testing                                                                    [2015-09-17 02:34:54,415] WARN Property topic is not valid (kafka.utils.VerifiableProperties)
    This is a Test for flume and Kafka intigration
    We can expect these messages into Flume SINK (HDFS)
    you can see on flume agent console window that flume is creating
    Files in HDFS

    produce_mesg
    You can see from Flume agent window that it would be storing those messages to the sink (HDFS):
    HDFS_SINK

    경축! 아무것도 안하여 에스천사게임즈가 새로운 모습으로 재오픈 하였습니다.
    어린이용이며, 설치가 필요없는 브라우저 게임입니다.
    https://s1004games.com

  2. Using Kafka as a SINK for Flume:
    Now we will pass messages from Netcat (Flume Source), which will go through Flume channel (In-Memory), finally getting Stored into Kafka and Can be accessed from Kafka Consumers.
    Create a new topic ‘testing3’:

     

    [student4@nbc-n4 flume]$ ~/kafka_2.10-0.8.2.0/bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic testing3

    Listing the Topics:

    [student4@nbc-n4 flume]$ ~/kafka_2.10-0.8.2.0/bin/kafka-topics.sh --list --zookeeper localhost:2181
    testing
    testing3

    Set-up Flume Conf:

    [student4@nbc-n4 flume]$ cat confFiles/exercise6.conf
    ##Configuring Components
    a1.sources=source1
    a1.channels=channel1
    a1.sinks=sink1
    ## Configuring Source
    a1.sources.source1.channels=channel1
    a1.sources.source1.type=netcat
    a1.sources.source1.bind=10.11.12.122
    a1.sources.source1.port=44444
    ## Configuring Channel
    a1.channels.channel1.type=memory
    a1.channels.channel1.capacity=10000
    a1.channels.channel1.transactionCapacity=1000
    ## Configuring Sink
    a1.sinks.sink1.type=org.apache.flume.sink.kafka.KafkaSink
    a1.sinks.sink1.topic=testing2
    a1.sinks.sink1.zookeeperConnect=localhost:2181
    a1.sinks.sink1.brokerList=localhost:9092
    a1.sinks.sink1.channel=channel1
    a1.sinks.sink1.batchSize=20

    Start Flume Agent:

    ./${flume_home}/bin/flume-ng agent -c ${flume_home}/conf/ -f ${flume_home/conf/exercise5.conf --name a1 -Dflume.root.logger=INFO,console

    Pass the Messages using Netcat and Access the same message via Kafka-consumers:

  3. [student4@nbc-n4 ~]$ # Passing Messages via Netcat
    [student4@nbc-n4 ~]$ nc 10.11.12.122 44444
    Hi
    OK
    This is Kafka-Flume Integration Testing
    OK
    We are using Kafka as Sink here
    OK
    You can access these message from Kafka-Consumers
    OK
    [student4@nbc-n4 ~]$ # Checking these messages in Kafka
    [student4@nbc-n4 ~]$ ./kafka_2.10-0.8.2.0/bin/kafka-console-consumer.sh --zookeeper localhost:2181 --topic testing3 --from-beginning
    Hi
    This is Kafka-Flume Integration Testing
    We are using Kafka as Sink here
    You can access these message from Kafka-Consumers


    3. Using Kafka as a CHANNEL for Flume:
    Now we will pass messages from Netcat (Flume Source), which will go through Kafka Topics (act as channel) and finally, can be accessed from Kafka Consumers.
    Below is the configuration file:

    [student4@nbc-n4 flume]$ cat confFiles/exercise7.conf
    #Name the components on this agent
    a1.sources = r1
    a1.channels = c1
    #Describe/configure the source
    a1.sources.r1.type = netcat
    a1.sources.r1.bind = localhost
    a1.sources.r1.port = 44444
    a1.channels.c1.type = org.apache.flume.channel.kafka.KafkaChannel
    a1.channels.c1.capacity = 10000
    a1.channels.c1.transactionCapacity = 1000
    a1.channels.c1.brokerList=localhost:9092
    a1.channels.c1.topic=testing3
    a1.channels.c1.zookeeperConnect=localhost:2181
    #Bind the source and sink to the channel
    a1.sources.r1.channels = c1

    Using above config file, you can try sending messages via Netcat and access them via Kafka-consumers.
     

본 웹사이트는 광고를 포함하고 있습니다.
광고 클릭에서 발생하는 수익금은 모두 웹사이트 서버의 유지 및 관리, 그리고 기술 콘텐츠 향상을 위해 쓰여집니다.
번호 제목 글쓴이 날짜 조회 수
공지 오라클 기본 샘플 데이터베이스 졸리운_곰 2014.01.02 86191
공지 [SQL컨셉] 서적 "SQL컨셉"의 샘플 데이타 베이스 SAMPLE DATABASE of ORACLE 가을의 곰을... 2013.02.10 78669
공지 [G_SQL] Sample Database 가을의 곰을... 2012.05.20 95405
22 Docker에서 SQL Server 컨테이너 이미지 구성 file 졸리운_곰 2020.01.23 1552
21 MSSQL 설치형 한글 환경으로 변경 file 졸리운_곰 2020.01.23 2885
20 PRIMARY KEY 와 FOREIGN KEY 를 전부 뽑아주는 쿼리 졸리운_곰 2018.12.16 1773
19 [MSSQL] CASE 문 . 조건에 따라 값 정하기 ! CASE WHEN THEN 졸리운_곰 2018.07.24 1565
18 Track Data Changes (SQL Server) file 졸리운_곰 2018.07.02 1302
17 Docker가 있는 SQL Server 2017 컨테이너 이미지를 실행 하는 빠른 시작 file 졸리운_곰 2018.06.26 1054
16 [MSSQL] Management Studio 이용해 데이터베이스 생성하기 file 졸리운_곰 2018.06.17 1281
15 [MSSQL - GROUP BY HAVING 을 이용한 중복 데이타 체크] file 졸리운_곰 2018.06.15 1275
14 [SQL] select 한 결과로 update 처리, SQL한문장, How to UPDATE from SELECT in SQL Server 졸리운_곰 2018.01.22 1470
13 UNION으로 결과 집합 조합 졸리운_곰 2017.08.27 1340
12 uniqueidentifier(Transact-SQL) file 가을의곰 2017.06.10 1754
11 하위 쿼리를 사용하여 다른 쿼리 또는 식에 쿼리 중첩 [MS-ACCESS : ms offce suit] 가을의곰 2017.06.10 1638
10 [MS-SQL] 테이블명, 컬럼명 검색 졸리운_곰 2017.04.17 1931
9 DB의 모든 테이블에서 데이터 검색 졸리운_곰 2017.04.17 1716
8 Microsoft SQL Server DBA 가이드-DBA라면 이정도는 알아야한다!!! file 졸리운_곰 2017.01.15 1302
7 SQL Server DBA 가이드 file 졸리운_곰 2017.01.15 1685
6 IDENTITY_INSERT가 OFF로 설정되면 ‘테이블명’ 테이블의 ID 열에 명시적 값을 삽입할 수 없습니다 file 졸리운_곰 2017.01.15 1502
5 MS SQL 서버에서 자동증가, autoincrement 처리 file 졸리운_곰 2017.01.15 1773
4 MS SQL 서버의 날짜, 시간 => 문자열 변환 포멧 설명 졸리운_곰 2017.01.15 1214
3 MS SQL 서버 코딩 표준 가이드 file 졸리운_곰 2017.01.14 1653
대표 김성준 주소 : 경기 용인 분당수지 U타워 등록번호 : 142-07-27414
통신판매업 신고 : 제2012-용인수지-0185호 출판업 신고 : 수지구청 제 123호 개인정보보호최고책임자 : 김성준 sjkim70@stechstar.com
대표전화 : 010-4589-2193 [fax] 02-6280-1294 COPYRIGHT(C) stechstar.com ALL RIGHTS RESERVED