레이블이 kafka인 게시물을 표시합니다. 모든 게시물 표시
레이블이 kafka인 게시물을 표시합니다. 모든 게시물 표시

2020년 8월 18일 화요일

kafka + logstash saving to hdfs (webhdfs)

상품추천 분석을 위한 로그들을 kafka broker에 꾸준히 쏴주고 있었다. 이제 이를 사용할 때가 되었다.

이 로그들을 하둡에 쌓을 필요가 생겼고 컨슈머를 구현할지, logstash의 webhdfs를 사용할지 고민하다가 logstash의 webhdfs output이 간단해서 이걸 사용하기로 하였다.


필요한 플러그인들이 설치되어있다고 가정하고 conf파일을 작성한다.

크게 input, filter, output 부분을 작성하면 된다. 필요한 부분을 검색해서 쓰긴했는데 나중에 시간날 때 제대로 사용법을 익혀두는게 좋을 것 같다. (꽤 편리한 것 같다.)


input {
     kafka {
            bootstrap_servers => "a:9092,b:9092,c:9092"
            topics => ["topic1","topic2"]
            decorate_events => true
            consumer_threads => 3
            group_id => "dev_recommend"
            }
}


모든 토픽의 파티션은 3개로 되어있기때문에 consumer_threads를 3으로 했다.

데이터를 가공하는 filter부분에서 토픽접근은 [@metadata][kafka][topic]로 하면되고 or조건은 or로 쓰면 되었다. c1은 unix time으로 된 컬럼이고 컬럼끼리의 구분자는 \t으로 되어있다면 아래처럼 작성해주면 되었다.


filter {
    if [@metadata][kafka][topic] == "topic1" or [@metadata][kafka][topic] == "topic2" {                                                                            
        dissect {
            mapping => {"message" => "%{col1} %{col2}   %{col3}   %{col4}   %{col5}   %{col6}   %{col7}   %{col8}   %{col9}   %{col10}  %{col11}"}
        }
        ruby {
                code => "
                      require 'time'
                      require 'date'

                      ts_new = event.get('[ts]').to_i/1000

                      event.set('day_ymd', Time.at(ts_new).localtime.strftime('%y-%m-%d'))
                      event.set('day_h', Time.at(ts_new).localtime.strftime('%H'))
                      event.set('day_ymdh', Time.at(ts_new).localtime.strftime('%y-%m-%d-%H'))
                  "
            }
    }

}

마지막으로 ouput 부분은 다음처럼 작성하면 된다.

output {

    if [@metadata][kafka][topic] == "topic1" {
        webhdfs {
            host => "hadoop url"
            port => 50070
            path => "/data/tsv/recommend-logs/%{day_ymd}/%{day_h}/topic1-log-%{day_ymdh}.log"
            user => "usr"
            codec => line { format => "%{message}"}
            }
        #stdout { codec => rubydebug }
    }
    else if [@metadata][kafka][topic] == "topic2" {
        webhdfs {
            host => "hadoop url"
            port => 50070
            path => "/data/tsv/recommend-logs/%{day_ymd}/%{day_h}/topic2-log-%{day_ymdh}.log"
            user => "usr"
            codec => line { format => "%{message}"}

            }
        #stdout { codec => rubydebug }
    }
}

이제 실행해서 화면에 찍어보면 다음과 같은 형식으로 로그가 찍힐것이다.


output file은 계속해서 append를 해야하기때문에 text 포맷이다. 이후 실행하면 yy-mm-dd/시간/파일 경로에 데이터가 잘 쌓이는 것을 확인할 수 있다. 
나중에 제대로 사용법(문법)을 익혀둬야겠다.

일단 작업을 돌려놓고 며칠간 사이즈 모니터링을 해봐야겠다.
nohup bin/logstash --path.data recom_data -f recom.conf > /dev/null 2>&1 &


2020년 8월 7일 금요일

Kafka + Spark Streaming saving to Cassandra

추천 시스템 서빙에서 메인 DB는 Cassandra를 사용하고 있는데 대부분이 data를 bulk insert 후 select만 해서 return 하는 구조이다. 하지만 최근 실시간 (방문)로그를 활용해야할 필요가 생겼고 이를 Cassandra에 적재하기로 하였다. 

만약 cassandra에 쌓는다면 real-time 환경에서 transaction에 문제가 없는지, 하루 쌓이는 양은 얼마 정도인지, compaction 전략은 어떤걸 써야할지, 피크 상황에서 서버가 버틸지 등의 고민을 하게 되었다.

이러한 시도가 성공할지, 실패할지는 먼저 구현을 해봐야 알 수 있기 때문에 일단 먼저 구현부터 하기로 했다. 만약 데이터양이 많은 것이 문제가 된다면 사이즈를 줄이는 방향으로 타협을 보는 것으로 문제를 해결해도 되는 상황이다.


구현 방법은 두 가지로 생각했다.
1. kafka -> logstash -> cassandra
2. kafka -> spark streaming -> cassandra

일단 1번은 logstash plugin이 존재하기는 하는데 4~5년전 자료라서 pass.
2번의 경우 어느 정도의 성능이 나올지 몰라서 궁금하기도 했고 2번으로 결정했다.


정리하면 real time으로 Kafka broker -> spark streaming(consum) -> cassandra 로 데이터를 흘러가는 것이 목표이다.

기존에 운영중인 kafka 데이터를 사용하기로 했다. kafka의 데이터는 map형식으로 들어오고 있다. 

{"a_key":"a_value","b_key":"b_value","c_key":"c_value","d_key":"d_value"}를 담을 class를 하나 만들었다.
case class TEST ( COL1:String, COL2:String, COL3:String, COL4:String )

이후 카프카와 연동을 해줬다.

val sparkConf = new SparkConf()
.setAppName("kapark2cassandra")
.set("spark.cassandra.auth.username", "cassandra")
.set("spark.cassandra.auth.password", "cassandra")
.set("spark.cassandra.connection.host", "cassandra cluster")

val ssc = new StreamingContext(sparkConf, Seconds(3))

val kafkaParams = Map[String, Object](
"bootstrap.servers" -> "zk cluster",
"key.deserializer" -> classOf[StringDeserializer],
"value.deserializer" -> classOf[StringDeserializer],
"group.id" -> "kapark2cassandra",
"auto.offset.reset" -> "earliest",
"enable.auto.commit" -> (false: java.lang.Boolean)
)
val topics = Array("kapark2cassandra")
val stream = KafkaUtils.createDirectStream[String, String](
ssc,
PreferConsistent,
Subscribe[String, String](topics, kafkaParams)
)

stream.foreachRDD { rdd =>
val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
stream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)
}

다음은 데이터 가공 및 카산드라에 전송이다.
ttl은 1주일(604800)으로 셋팅했다. 

stream.foreachRDD { rdd =>
val spark = SparkSessionSingleton.getInstance(rdd.sparkContext.getConf)
import spark.implicits._
val ods = rdd
.map(record => record.value.toString.drop(1).dropRight(1).split(",(?=([^\"]*\"[^\"]*\")*[^\"]*$)", -1)
.map(_.split(":(?=([^\\\"]*\\\"[^\\\"]*\\\")*[^\\\"]*$)")).map(arr => (arr(0) -> arr(1))).toMap)
.map(rec =>
TEST(
rec.get("\"COL1\"").mkString.replace("\"", "").trim
, rec.get("\"COL2\"").mkString.replace("\"", "").trim
, rec.get("\"COL3\"").mkString.replace("\"", "").trim
, rec.get("\"COL4\"").mkString.replace("\"", "").trim
)
).toDF()

ods.createOrReplaceTempView("kapark2cassandra")

var result = spark.sql(
"""
select from_unixtime(col1/1000,"yyyy-MM-dd HH:mm:ss") as col1
, col2
, col3
, col4
from kapark2cassandra
where col2 = some
""")
result.write.format("org.apache.spark.sql.cassandra")
.options(Map("table" -> "test" , "keyspace" -> "real_time", "ttl" -> "604800"))
.mode(org.apache.spark.sql.SaveMode.Append).save()
}

ssc.start()
ssc.awaitTermination()
저기서 다른 데이터 소스와 join, group by해서 사용할 수 있으면 그렇게 하려고 하는데 스트리밍이 밀리지 않는 선에서 처리가 가능한 양이어야 할 것이다.
object SparkSessionSingleton {
@transient private var instance: SparkSession = _
def getInstance(sparkConf: SparkConf): SparkSession = {
if (instance == null) {
instance = SparkSession
.builder
.config(sparkConf)
.getOrCreate()
}
instance
}
}
마지막으로 val spark = SparkSessionSingleton.getInstance(rdd.sparkContext.getConf) 부분을 만들기 위한 ojbect 생성해주고 
def main(args: Array[String]) {
new kapark2cassandra().run(args)
}
수행했더니 데이터가 카산드라에 잘 들어간다.

통계를 보도록 하자.
대충 1큐에 5천건~8천건 정도이다.
3초로 했을 때 조금씩 튀는 부분이 있는데 아무래도 3초는 무리인것같다.

cassandra 지표를 봐도 딱히 튀는 부분도 없고 결국 final result는 약 100 rows씩 들어가기 때문에 부하 테스트를 하기에는 작은 양이었다.

결론
real time data를 cassandra에 넣으려면 kafka + spark streaming으로 처리하는 것도 하나의 옵션이다.





2020년 5월 26일 화요일

Apache Kafka Cluster 설치 및 테스트

블로그 포스팅하려고 집 컴퓨터에 카프카를 세팅하려고했는데 확실히 집에서 하려고 하니 제약사항이 많았다. 그래도 이왕 해보기로 한거 끝까지 해본다.

일단 생각하는 것은 kafka broker는 3대를 클러스터링하고 컨슈머,프로듀서를 각각 세팅해놓으려고 한다.

먼저 centos7 가상환경을 3대를 준비했다. (나중에 프로듀서, 컨슈머로 사용하려고 2대 더 준비했다. 총5대 )

/etc/hosts에 kafka 클러스터로 쓸 host 3대를 추가한다.

192.168.20.130 kafka-srv01
192.168.20.131 kafka-srv02
192.168.20.128 kafka-srv03

먼저 java를 설치한다. (zulu 1.8)



kafka cluster는 zookeeper가 관리하기 때문에 zookeeper와 함께 설치되어야한다.
현재 나온 카프카 최신버전은 2.5.0 버전이고 카프카에 포함되어있는 주키퍼가 있지만 주키퍼를 별도로 설치했다. (주키퍼 3.5.8)

디렉토리 구조는 다음처럼 가져갔다.



먼저 주키퍼 부터 세팅하자.

먼저 디렉토리를 만들었다.
mkdir -p /data01/tmp/zookeeper

각 노드에서 id를 세팅해주었다.

echo 1 > /data01/tmp/zookeeper/myid
echo 2 > /data01/tmp/zookeeper/myid
echo 3 > /data01/tmp/zookeeper/myid

다음으로 /data01/sw/zookeeper/conf에서 zoo_sample.cfg를 복사해서 zoo.cfg를 만들었다.

아래처럼 수정한다.

dataDir=/data01/tmp/zookeeper
(생략)
initLimit=5
syncLimit=2

server.1=kafka-srv01:2888:3888
server.2=kafka-srv02:2888:3888
server.3=kafka-srv03:2888:3888


server.myid=호스트:통신용포트1:통신용포트2만 작성해주면 되서 쉽다.
myid는 주키퍼 클러스터에서 각 서버에 부여하는 고유한 서버 번호이다. 따라서 포스팅 내용처럼 tmp 디렉토리가 아닌 다른 디렉토리를 만들어서 보관하자.


다음은 카프카 설정파일을 수정한다. (vi server.properties)

mkdir -p /data01/tmp/kafka-logs

broker.id=1
log.dirs=/data01/tmp/kafka-logs
listeners=PLAINTEXT://:9092
advertised.listeners=PLAINTEXT://kafka-srv01:9092
zookeeper.connect=kafka-srv01:2181,kafka-srv02:2181,kafka-srv03:2181

broker.id=2
log.dirs=/data01/tmp/kafka-logs
listeners=PLAINTEXT://:9092
advertised.listeners=PLAINTEXT://kafka-srv02:9092
zookeeper.connect=kafka-srv01:2181,kafka-srv02:2181,kafka-srv03:2181

broker.id=3
log.dirs=/data01/tmp/kafka-logs
listeners=PLAINTEXT://:9092
advertised.listeners=PLAINTEXT://kafka-srv03:9092
zookeeper.connect=kafka-srv01:2181,kafka-srv02:2181,kafka-srv03:2181


브로커id는 안주면 카프카가 자동으로 세팅해주고 굳이 주키퍼랑 맞출 필요도 없긴 하지만 이쁘게 통일하기 위해 주키퍼-카프카 id를 맞췄다.
그리고 zookeeper.connect는 브로커가 주키퍼에 접속할 때의 접속정보다.


설정은 끝났다.

이제 구동해보자. 주키퍼를 먼저 띄우고 그다음 카프카(브로커)를 띄워야한다. 각 노드의 순서는 상관 없다. 반대로 서비스를 내릴 때에는 카프카(브로커)를 내린 후에 주키퍼를 내린다.

주키퍼 실행
zkServer.sh start
카프카 실행
bin/kafka-server-start.sh -daemon config/server.properties

카프카 서버 로그를 보니 잘 올라온 것을 확인할 수 있었다.



이제 토픽을 만들고 테스트를 해보자.

test-topic을 만들었다. 파티션은 3을 주고 레플리카도 3을 줬다.
./kafka-topics.sh --zookeeper kafka-srv01:2181,kafka-srv02:2181,kafka-srv03:2181 --create --topic test-topic --partitions 3 --replication-factor 3

잘 만들어졌는지 확인해본다.
./kafka-topics.sh --zookeeper kafka-srv01:2181,kafka-srv02:2181,kafka-srv03:2181 --describe --topic test-topic



describe만 살펴보자.


Leader는 각 파티션의 현재 Leader 복제본이 어떤 브로커에 있는지 알려준다. Replicas는 각 파티션의 복제본을 보유하고 있는 브로커의 리스트이다.

Isr은 In-Sync Replicas의 약자로 복제본 중에서 Leader Replica랑 동기화가 되고 있는 복제본을 소유하고 있는 브로커의 리스트이다. 브로커가 장애가 나거나 동기화가 안됐을 때 Isr에 포함되지 않는다.

테스트 토픽에 잘 던지고 받는지 console로 확인해보자.
나중에 이것저것 더 테스트해보려고 별도로 두대의 서버를 만들고 카프카를 설치했다.

각각 서버에서 프로듀서와 컨슈머를 띄웠다.
./kafka-console-producer.sh --broker-list kafka-srv01:9092,kafka-srv02:9092,kafka-srv03:9092 --topic test-topic
./kafka-console-consumer.sh --bootstrap-server kafka-srv01:9092,kafka-srv02:9092,kafka-srv03:9092 --topic test-topic


프로듀서에서 hello world, park, su, seong, parksuseong을 던지면 컨슈머에서 잘 받는 것을 확인할 수 있다.

kafka-producer


kafka-consumer



사실 카프카 혼자서는 서비스를 하기엔 무리이고 카프카를 이용한 생태계를 구축해야하는데 시간이 날지 모르겠다. 그래도 추후 fluentd나 logstash, elk, spark streaming 정도까지는 포스팅해보고싶은데 일단 노력해봐야지.. 집 컴퓨터라서 혼자서 세팅하는데 엄청난 인내심이 필요하기 때문이다.

2022년 회고

 올해는 블로그 포스팅을 열심히 못했다. 개인적으로 지금까지 경험했던 내용들을 리마인드하자는 마인드로 한해를 보낸 것 같다.  대부분의 시간을 MLOps pipeline 구축하고 대부분을 최적화 하는데 시간을 많이 할애했다. 결국에는 MLops도 데이...