이 문서에서는 대규모 데이터 분석을 위해 Apache Impala를 활용하여 SQL 기반 데이터 웨어하우징 작업을 수행하는 방법과, Apache Flink를 사용하여 Apache Kafka와 연동하여 실시간 스트림 데이터를 처리하는 애플리케이션을 개발하는 과정을 다룹니다.
빅데이터 클러스터 환경 구성
Impala와 Flink를 사용하기 위해서는 Hadoop, Impala, Kafka, Flink, Zookeeper, Ranger 등의 주요 서비스가 설치된 빅데이터 클러스터 환경이 필요합니다. 클러스터 구성 시 다음과 같은 사항을 고려해야 합니다.
- 서비스 선택: Impala, Flink, Kafka, Zookeeper는 필수적으로 설치해야 합니다.
- 노드 구성: 데이터 노드(DN), 노드 매니저(NM), 블록(B) 등 필요한 역할을 수행하는 서비스가 각 노드에 적절히 배포되어야 합니다. 특히 Impalad 서비스는 데이터 처리 노드에 활성화되어야 합니다.
- 네트워크 설정: 외부에서 클러스터에 접근할 수 있도록 네트워크 보안 그룹 또는 방화벽 규칙을 구성하여 SSH 및 서비스 포트 접근을 허용해야 합니다.
클러스터 준비가 완료되면 SSH 클라이언트를 사용하여 클러스터 마스터 노드에 접속합니다.
Impala를 활용한 데이터 관리
Impala는 Hadoop 분산 파일 시스템(HDFS)에 저장된 데이터를 위한 대규모 병렬 처리(MPP) SQL 쿼리 엔진입니다. C++와 Java로 개발되었으며, 다른 Hadoop 기반 SQL 엔진에 비해 높은 성능과 낮은 지연 시간을 제공합니다.
사용자 관리 시나리오
가상의 고객 관리 시스템에서 사용자 정보를 관리하는 애플리케이션을 개발한다고 가정합니다. Impala 클라이언트를 사용하여 다음과 같은 데이터베이스 작업을 수행합니다.
- 사용자 상세 정보 테이블 생성
- 테이블에 학력, 직책 정보 추가
- 특정 사용자 ID로 이름과 주소 조회
- 모든 작업 완료 후 사용자 정보 테이블 삭제
Impala 클라이언트 설정
Impala 작업을 수행하기 위해서는 클러스터에 Impala 클라이언트를 설치해야 합니다.
- 클러스터 관리 페이지에서 Impala 클라이언트 패키지를 다운로드합니다.
- 다운로드한 아카이브 파일을 서버의 임시 디렉터리(예:
/tmp/client_packages)에 업로드하고 압축을 해제합니다.cd /tmp/client_packages tar -vxf Impala_Client.tar tar -vxf Impala_ClientConfig.tar - 설치 스크립트를 실행하여 클라이언트를 특정 디렉터리(예:
/opt/bigdata/client)에 설치합니다.cd Impala_ClientConfig ./install.sh /opt/bigdata/client - 설치 디렉터리로 이동하여 환경 변수를 설정하고 Impala 셸을 실행합니다.
cd /opt/bigdata/client source bigdata_env impala-shell - Impala 셸을 종료하려면
quit;명령을 사용합니다.
Impala SQL 예제
이제 Impala 셸에서 SQL 명령을 실행하여 사용자 정보를 관리합니다.
내부(Managed) 테이블 작업
내부 테이블은 Impala가 HDFS 내의 데이터와 메타데이터를 모두 관리하는 테이블입니다.
-- 1. 사용자 상세 정보 테이블 생성
CREATE TABLE customer_details (
customer_id STRING,
full_name STRING,
gender STRING,
age INT,
city STRING
);
-- 2. 사용자 정보 추가
INSERT INTO customer_details (customer_id, full_name, gender, age, city) VALUES
("CUST001", "김철수", "남", 29, "서울");
-- 3. 테이블에 학력 및 직책 정보 컬럼 추가
ALTER TABLE customer_details ADD COLUMNS (
education_level STRING,
job_title STRING
);
-- 4. 특정 사용자 ID로 이름과 도시 조회
SELECT full_name, city FROM customer_details WHERE customer_id='CUST001';
-- 5. 사용자 정보 테이블 삭제
DROP TABLE customer_details;
외부(External) 테이블 작업
외부 테이블은 데이터 파일을 Impala 외부에서 관리하며, Impala는 메타데이터만 관리합니다. DROP TABLE 시 데이터 파일은 삭제되지 않습니다.
-- 1. 외부 테이블 생성 (등록 연도로 파티션)
CREATE EXTERNAL TABLE customer_details_ext (
customer_id STRING,
full_name STRING,
gender STRING,
age INT,
city STRING
)
PARTITIONED BY (register_year STRING)
ROW FORMAT DELIMITED FIELDS TERMINATED BY ','
LINES TERMINATED BY '\n'
STORED AS TEXTFILE
LOCATION '/hive/customer_data_ext';
-- 2. 데이터를 파티션에 삽입
INSERT INTO customer_details_ext PARTITION(register_year='2023') VALUES
("CUST002", "이영희", "여", 25, "부산");
-- 3. 삽입된 데이터 확인
SELECT * FROM customer_details_ext;
파일 데이터 로드
외부 파일에 저장된 데이터를 Impala 테이블로 로드할 수 있습니다.
- 샘플 데이터 파일 생성 (
customer_data.txt)CUST003,박민수,남,33,대구 CUST004,최지영,여,27,인천 - HDFS에 데이터 파일 업로드
hdfs dfs -put customer_data.txt /data/customer_info - Impala 셸에서 데이터 로드 명령 실행
impala-shell LOAD DATA INPATH '/data/customer_info/customer_data.txt' INTO TABLE customer_details_ext PARTITION (register_year='2023'); SELECT * FROM customer_details_ext; - 작업 완료 후 테이블 삭제
DROP TABLE customer_details_ext;
Flink를 활용한 실시간 스트림 처리
이 섹션에서는 Apache Flink와 Kafka를 연동하여 실시간으로 수신되는 메시지에 특정 접두사를 추가하여 처리하고 출력하는 간단한 스트리밍 애플리케이션을 개발합니다.
실시간 스트리밍 시나리오
초당 하나의 메시지가 Kafka 토픽으로 전송된다고 가정합니다. Flink 애플리케이션은 이 메시지를 실시간으로 소비하여, 메시지 내용 앞에 특정 접두사(예: "[Flink Message]")를 붙여 다시 Kafka 토픽으로 발행합니다.
JDK 설치
Flink 애플리케이션 개발 및 실행을 위해 Java Development Kit 8 (JDK 8)이 필요합니다.
wget https://download.java.net/java/GA/jdk8u341/8/GPL/openjdk-8u341-b08_linux-x64_bin.tar.gz
tar -zxvf openjdk-8u341-b08_linux-x64_bin.tar.gz
다운로드 후 압축을 해제하고, 시스템의 JAVA_HOME 환경 변수 및 PATH에 JDK 경로를 설정합니다.
Flink Maven 프로젝트 생성
IDE (예: Eclipse)에서 새로운 Maven 프로젝트를 생성합니다.
- Group Id:
com.example - Artifact Id:
flink-kafka-processor
IDE에서 JDK 경로 설정
프로젝트가 올바른 JDK 버전(JDK 8)을 사용하도록 IDE에서 프로젝트의 JRE 시스템 라이브러리 설정을 확인하고 필요에 따라 변경합니다. 기존에 설정된 다른 JRE 버전이 있다면 제거하고 새로 설치한 JDK 8을 Workspace default JRE로 설정합니다.
POM 파일 설정
pom.xml 파일을 열어 Flink 및 Kafka 커넥터 의존성을 추가하고 빌드 설정을 구성합니다.
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>flink-kafka-processor</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<flink.version>1.13.0</flink.version>
<java.version>1.8</java.version>
<maven.compiler.source>${java.version}</maven.compiler.source>
<maven.compiler.target>${java.version}</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java_2.11</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka_2.11</artifactId>
<version>${flink.version}</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.8.1</version>
<configuration>
<source>${java.version}</source>
<target>${java.version}</target>
</configuration>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.2.4</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<configuration>
<artifactSet>
<excludes>
<exclude>org.apache.flink:force-shading</exclude>
<exclude>com.google.code.findbugs:jsr305</exclude>
<exclude>org.slf4j:*</exclude>
<exclude>log4j:*</exclude>
</excludes>
</artifactSet>
<filters>
<filter>
<artifact>*:*</artifact>
<excludes>
<exclude>META-INF/*.SF</exclude>
<exclude>META-INF/*.DSA</exclude>
<exclude>META-INF/*.RSA</exclude>
</excludes>
</filter>
</filters>
<transformers>
<transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
<mainClass>com.example.flink.KafkaMessageGenerator</mainClass>
</transformer>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
POM 파일 저장 후 Maven이 의존성을 다운로드하도록 기다립니다.
애플리케이션 개발
src/main/java 폴더 아래에 com.example.flink 패키지를 생성하고, 그 안에 KafkaMessageGenerator.java 클래스를 생성합니다. 이 클래스는 Kafka로 메시지를 발행하는 Flink Producer 역할을 합니다.
package com.example.flink;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.source.SourceFunction;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
import org.apache.flink.streaming.util.serialization.SimpleStringSchema;
import java.util.Properties;
public class KafkaMessageGenerator {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1); // 단일 태스크로 실행
ParameterTool params = ParameterTool.fromArgs(args);
// Kafka 프로듀서 설정을 위한 Properties 객체 생성
Properties kafkaProps = params.getProperties();
// 실시간으로 메시지를 생성하는 소스 추가
DataStream<String> messageStream = env.addSource(new IncrementalMessageSource());
// Kafka 싱크 (Producer) 설정 및 추가
// 토픽 이름과 Kafka 브로커 주소는 실행 시 인자로 받음
messageStream.addSink(
new FlinkKafkaProducer<String>(
params.getRequired("output.topic"), // 메시지를 보낼 Kafka 토픽
new SimpleStringSchema(), // 메시지를 직렬화하는 스키마
kafkaProps // Kafka 프로듀서 속성
)
);
// Flink 애플리케이션 실행
env.execute("Flink Kafka Message Generator");
}
/**
* 일정 간격으로 순차적인 문자열 메시지를 생성하는 소스 함수.
*/
public static class IncrementalMessageSource implements SourceFunction<String> {
private static final long serialVersionUID = 1L; // 직렬화 ID
private volatile boolean isRunning = true;
private long messageCounter = 0;
@Override
public void run(SourceContext<String> ctx) throws Exception {
while (isRunning) {
// 특정 접두사를 포함한 메시지 생성
ctx.collect("[Flink Message] ID: " + (messageCounter++));
Thread.sleep(1000); // 1초 간격으로 메시지 생성
}
}
@Override
public void cancel() {
isRunning = false;
}
}
}
애플리케이션 패키징 및 배포
Maven을 사용하여 Flink 애플리케이션을 JAR 파일로 패키징합니다.
- 프로젝트 루트에서 다음 Maven 명령을 실행하여 JAR 파일을 생성합니다.
mvn clean package성공적으로 빌드되면
target/디렉터리에flink-kafka-processor-1.0-SNAPSHOT.jar와 같은 이름의 JAR 파일이 생성됩니다. - 생성된 JAR 파일을 클러스터의 특정 디렉터리(예:
/home/user/flink_jobs)로 복사합니다.scp /path/to/your/project/target/flink-kafka-processor-1.0-SNAPSHOT.jar user@<클러스터_IP>:/home/user/flink_jobs
Flink 클라이언트 설치
Impala 클라이언트와 유사하게, Flink 애플리케이션을 실행하기 위해 클러스터에 Flink 클라이언트를 설치합니다.
- 클러스터 관리 페이지에서 Flink 클라이언트 패키지를 다운로드하고 서버에 업로드합니다.
- 압축을 해제하고 설치 스크립트를 실행하여 Flink 클라이언트를 설치합니다.
cd /tmp/client_packages tar -vxf Flink_Client.tar tar -vxf Flink_ClientConfig.tar cd Flink_ClientConfig ./install.sh /opt/bigdata/flink_client - 설치 디렉터리로 이동하여 환경 변수를 설정합니다.
cd /opt/bigdata/flink_client source bigdata_env
Kafka 환경 준비 및 Flink 애플리케이션 실행
Flink 애플리케이션을 실행하기 전에 Kafka 토픽을 생성하고, Flink 애플리케이션의 출력을 확인할 Kafka 소비자(Consumer)를 준비합니다.
Kafka Broker 및 Zookeeper 정보 확인
클러스터의 Kafka Broker IP 주소와 Zookeeper 주소를 확인합니다. 이는 Kafka 클라이언트 명령 및 Flink 애플리케이션 실행 시 필요합니다.
Kafka 토픽 생성
메시지를 주고받을 Kafka 토픽을 생성합니다. 여기서는 flink-input-topic이라는 토픽을 사용합니다.
- Kafka 클라이언트 디렉터리로 이동합니다.
cd /opt/Bigdata/components/<Hadoop_Version>/Kafka/client/install_files/kafka/bin - 스크립트 실행 권한을 부여합니다.
chmod 777 kafka-topics.sh kafka-run-class.sh kafka-console-producer.sh kafka-console-consumer.sh - 토픽을 생성합니다.
<Zookeeper_IP>는 실제 Zookeeper 서버 IP로 대체합니다../kafka-topics.sh --create --zookeeper <Zookeeper_IP>:2181/kafka --topic flink-input-topic --replication-factor 2 --partitions 2 - 생성된 토픽 목록을 확인합니다.
./kafka-topics.sh --list --zookeeper <Zookeeper_IP>:2181/kafka
Kafka 컨슈머 시작
새로운 터미널 창을 열고 Kafka 컨슈머를 시작하여 Flink 애플리케이션이 발행할 메시지를 실시간으로 확인합니다. <Kafka_Broker_IP>는 실제 Kafka Broker IP로 대체합니다.
ssh user@<클러스터_IP>
cd /opt/Bigdata/components/<Hadoop_Version>/Kafka/client/install_files/kafka/bin
./kafka-console-consumer.sh --topic flink-input-topic --bootstrap-server <Kafka_Broker_IP>:9092 --from-beginning
Flink 애플리케이션(Producer) 실행
또 다른 터미널 창을 열고 Flink 애플리케이션을 실행합니다. Flink 애플리케이션은 일반적으로 root가 아닌 omm 또는 flink와 같은 특정 사용자 계정으로 실행됩니다.
- 클러스터에 SSH로 접속 후, Flink 클라이언트 설치 경로로 이동하여 환경 변수를 로드합니다.
ssh user@<클러스터_IP> su omm # 또는 Flink 실행 권한이 있는 사용자 cd /opt/bigdata/flink_client source bigdata_env FlinkYarnSessionCli프로세스가 실행 중인지 확인합니다.jps만약
FlinkYarnSessionCli프로세스가 보이지 않는다면, Yarn 세션을 시작해야 합니다.cd Flink/flink/bin ./yarn-session.sh -d & # 백그라운드에서 Yarn 세션 시작- Flink 애플리케이션을 실행합니다.
<Kafka_Broker_IP>는 실제 Kafka Broker IP로 대체합니다.flink run \ --class com.example.flink.KafkaMessageGenerator \ /home/user/flink_jobs/flink-kafka-processor-1.0-SNAPSHOT.jar \ --output.topic flink-input-topic \ --bootstrap.servers <Kafka_Broker_IP>:9092이제 Kafka 컨슈머 터미널에서 Flink 애플리케이션이 Kafka로 발행하는 메시지를 실시간으로 확인할 수 있습니다. 메시지는 "[Flink Message] ID: N"과 같은 형식으로 1초마다 출력될 것입니다.