스트리밍 데이터 수집
스트리밍 데이터 수집 (Streaming Data Ingest)
전통적으로 Hive에 새 데이터를 추가하려면 많은 데이터를 HDFS에 모은 뒤 주기적으로 새 파티션을 추가해야 했어요 — 본질적으로 "배치 삽입"이죠. 기존 파티션에 새 데이터를 삽입하는 것은 허용되지 않았어요. Hive Streaming API는 데이터를 Hive에 연속적으로 흘려보낼 수 있게 해주며, 들어오는 데이터를 작은 레코드 배치로 기존 파티션/테이블에 계속 커밋할 수 있게 해줘요.
출처: 문서
본문
Hive 3 스트리밍 API
Hive 3 Streaming API Documentation - Hive 3에서 사용 가능한 새 API
Hive HCatalog 스트리밍 API
전통적으로 Hive에 새 데이터를 추가하려면 많은 양의 데이터를 HDFS에 모은 뒤 주기적으로 새 파티션을 추가해야 했어요. 이는 본질적으로 "배치 삽입(batch insertion)"이에요. 기존 파티션에 새 데이터를 삽입하는 것은 허용되지 않아요. Hive Streaming API는 데이터를 Hive에 연속적으로 흘려보낼 수 있게 해줘요. 들어오는 데이터를 소량의 레코드 배치로 기존 Hive 파티션 또는 테이블에 계속 커밋할 수 있어요. 데이터가 커밋되면 이후에 시작된 모든 Hive 쿼리에서 즉시 볼 수 있게 돼요.
이 API는 지속적으로 데이터를 생성하는 Flume과 Storm 같은 스트리밍 클라이언트를 대상으로 해요. 스트리밍 지원은 Hive의 ACID 기반 insert/update 지원 위에 구축돼 있어요 (참고: Hive Transactions).
Hive 스트리밍 API의 클래스와 인터페이스는 크게 두 집합으로 분류돼요. 첫 번째 집합은 커넥션과 트랜잭션 관리를 지원하고, 두 번째 집합은 I/O 지원을 제공해요. 트랜잭션은 metastore가 관리해요. 쓰기는 HDFS에 직접 수행돼요.
파티션되지 않은(unpartitioned) 테이블에 대한 스트리밍도 지원돼요. API는 Hive 0.14부터 Kerberos 인증을 지원해요.
패키징에 대한 참고: API는 Java 패키지 org.apache.hive.hcatalog.streaming에 정의되어 있으며, Hive의 hive-hcatalog-streaming Maven 모듈의 일부예요.
스트리밍 변경 API (Streaming Mutation API)
릴리스 2.0.0부터 Hive는 Hive의 ACID 기능을 사용해 transactional 테이블에 레코드를 변경(insert/update/delete)하는 또 다른 API를 제공해요. 자세한 내용과 이 문서에서 설명하는 스트리밍 데이터 수집 API와의 비교는 HCatalog Streaming Mutation API를 참고해요.
스트리밍 요구 사항 (Streaming Requirements)
스트리밍을 사용하려면 몇 가지가 필요해요.
- 스트리밍을 위해 ACID 지원을 활성화하려면 hive-site.xml에 다음 설정이 필요해요:
- 테이블 생성 시 "stored as orc" 를 지정해야 해요. 현재 ORC 저장 포맷만 지원돼요.
- 생성 시 테이블에 tblproperties("transactional"="true")가 설정되어야 해요.
- Hive 테이블은 버킷화(bucketed)되어야 하지만 정렬되지는 않아야 해요. 즉 생성 시 "clustered by (colName) into 10 buckets" 같은 것을 지정해야 해요. 버킷 수는 이상적으로 스트리밍 작성자 수와 같아야 해요.
- 스트리밍 클라이언트 프로세스의 사용자는 테이블/파티션에 쓰고 테이블에 파티션을 만들 수 있는 필요한 권한이 있어야 해요.
- (임시 요구 사항) 스트리밍 테이블에 쿼리를 실행할 때 클라이언트는 다음을 설정해야 해요
- hive.vectorized.execution.enabled 를 false로 (Hive 0.14.0 미만 버전용)
- hive.input.format 를 org.apache.hadoop.hive.ql.io.HiveInputFormat으로
제약 사항 (Limitations)
현재 스트리밍 API는 기본적으로 구분(delimited) 입력 데이터(CSV, 탭 구분 등)와 JSON(엄격한 문법) 형식 데이터만 스트리밍할 수 있게 지원해요. 다른 입력 포맷 지원은 RecordWriter 인터페이스의 추가 구현으로 제공할 수 있어요.
현재 대상 테이블의 포맷으로는 ORC만 지원돼요.
API 사용 (API Usage)
트랜잭션 및 커넥션 관리
HiveEndPoint
HiveEndPoint 클래스는 연결할 Hive 엔드 포인트를 설명해요. 데이터베이스, 테이블, 파티션 이름을 설명해요. 그것에 newConnection 메서드를 호출하면 스트리밍 목적의 Hive MetaStore에 연결을 수립해요. StreamingConnection 객체를 반환해요. 같은 엔드포인트에 여러 연결을 수립할 수 있어요. 그런 다음 StreamingConnection을 사용해 I/O를 수행할 새 트랜잭션을 시작할 수 있어요.
데이터가 지속적으로 스트리밍되는 구성에서는 정기적으로 새 파티션에 데이터가 추가될 가능성이 높아요. Hive 관리자가 필요한 파티션을 미리 만들거나, 스트리밍 클라이언트가 필요할 때 만들 수 있어요. HiveEndPoint.newConnection()은 파티션을 자동 생성할지 여부를 나타내는 boolean 인자를 받아요. 파티션 생성은 원자적 동작이므로 여러 클라이언트가 파티션을 만들려 경쟁할 수 있지만 하나만 성공하므로, 스트리밍 클라이언트는 파티션을 만들 때 동기화할 필요가 없어요.
트랜잭션은 전통적인 데이터베이스 시스템과 약간 다르게 구현돼요. 각 트랜잭션에는 id가 있고, 여러 트랜잭션이 "Transaction Batch"로 그룹화돼요. 이는 여러 트랜잭션의 레코드를 더 적은 파일(트랜잭션당 1개 파일 대신)로 묶는 데 도움을 줘요. 연결 후 스트리밍 클라이언트는 먼저 새 트랜잭션 배치를 요청해요. 응답으로 트랜잭션 배치의 일부인 트랜잭션 ID 집합을 받아요. 이후 클라이언트는 새 트랜잭션을 시작해 한 번에 하나의 트랜잭션 id를 소비해요. 클라이언트는 트랜잭션당 하나 이상의 레코드를 write()하고, 다음 트랜잭션으로 전환하기 전에 현재 트랜잭션을 commit하거나 abort해요. 각 TransactionBatch.write() 호출은 I/O 시도를 자동으로 현재 Txn ID와 연결해요. 스트리밍 클라이언트 프로세스의 사용자는 파티션 또는 테이블에 대한 쓰기 권한이 필요해요. 특정 사용자로 연결을 획득하려면 Kerberos 기반 인증이 필요해요. 아래 보안 스트리밍 예제를 참고해요.
동시성 참고: 여러 TransactionBatch에서 동시에 I/O를 수행할 수 있어요. 하지만 트랜잭션 배치 내의 트랜잭션은 순차적으로 소비되어야 해요.
자세한 내용은 HiveEndPoint Javadoc을 참고해요. 일반적으로 사용자는 HiveEndPoint 객체로 대상 정보를 설정한 뒤 newConnection을 호출해 연결을 만들고 StreamingConnection 객체를 받아요.
StreamingConnection
StreamingConnection 클래스는 트랜잭션 배치를 획득하는 데 사용돼요. HiveEndPoint가 연결을 제공하면 애플리케이션은 일반적으로 fetchTransactionBatch를 호출하고 일련의 트랜잭션을 쓰는 루프에 들어가요. 종료할 때 애플리케이션은 close를 호출해야 해요. 자세한 내용은 Javadoc을 참고해요.
TransactionBatch
TransactionBatch는 일련의 트랜잭션을 쓰는 데 사용돼요. 각 버킷에서 TxnBatch마다 HDFS에 파일이 하나 생성돼요. API는 각 레코드를 검사해 어떤 버킷에 속하는지 결정하고 적절한 버킷에 써요. 테이블에 버킷이 5개 있으면 (컴팩션이 시작되기 전에) TxnBatch에 5개 파일(일부는 비어 있을 수 있음)이 있어요. Hive 1.3.0 이전에는 API의 버킷 계산 로직 결함으로 레코드가 버킷에 잘못 분배되어, 버킷 조인 알고리즘을 사용하는 쿼리에서 잘못된 데이터가 반환될 수 있었어요.
TxnBatch의 각 트랜잭션에 대해 애플리케이션은 beginNextTransaction, write, 그리고 적절히 commit 또는 abort를 호출해요. 자세한 내용은 Javadoc을 참고해요. 한 트랜잭션에는 둘 이상의 파티션 데이터가 포함될 수 없어요.
TransactionBatch의 트랜잭션은 hive.txn.timeout초 후 커밋 또는 중단되지 않으면 결국 metastore에 의해 만료돼요. TransactionBatch 클래스는 배치에서 사용되지 않는 트랜잭션의 수명을 연장하는 heartbeat() 메서드를 제공해요. 좋은 경험 법칙은 TransactionBatch를 만든 후 (hive.txn.timeout/2) 간격으로 heartbeat()를 호출하는 것이에요. 비활성 트랜잭션을 유지하기에 충분하면서도 metastore를 불필요하게 부하하지 않아요.
사용 지침 (Usage Guidelines)
일반적으로 각 트랜잭션에 포함되는 이벤트가 많을수록 더 많은 처리량을 달성할 수 있어요. 특정 수의 이벤트 후 또는 특정 시간 간격 후에 커밋하는 것이 일반적이며, 둘 중 먼저 오는 것 기준으로 해요. 후자는 이벤트 흐름 속도가 가변적일 때 트랜잭션이 너무 오래 열려 있지 않도록 보장해요. 단일 트랜잭션에 포함할 수 있는 데이터 양에는 실질적인 제한이 없어요. 유일한 우려는 트랜잭션이 실패할 경우 다시 재생(replay)해야 할 데이터 양이에요. TransactionBatch 개념은 StreamingAPI가 HDFS에 생성하는 파일 수를 줄이는 역할을 해요. 주어진 배치의 모든 트랜잭션은 (버킷당) 같은 물리적 파일에 쓰므로, 파티션은 열린 트랜잭션이 있는 배치의 가장 이른 트랜잭션 수준까지만 컴팩션될 수 있어요. 따라서 TransactionBatch를 지나치게 크게 만들면 안 돼요. 일부 시간 후 TransactionBatch를 닫는 타이머를 포함하는 것이 합리적이에요(사용되지 않은 트랜잭션이 있어도).
참고: Hive 1.3.0부터 TxnBatch.close()를 호출하면 현재 TxnBatch의 모든 미사용 트랜잭션이 중단(abort)돼요.
HiveConf 객체에 대한 참고
HiveEndPoint.newConnection()은 HiveConf 인자를 받아요. null로 설정하거나 미리 만든 HiveConf 객체를 제공할 수 있어요. null이면 내부적으로 HiveConf 객체를 만들어 연결에 사용해요. HiveConf 객체가 인스턴스화될 때 hive-site.xml이 들어 있는 디렉토리가 Java 클래스패스의 일부라면 HiveConf 객체는 그 값으로 초기화돼요. hive-site.xml이 없으면 객체는 기본값으로 초기화돼요. 이 객체를 미리 만들고 여러 연결에서 재사용하면, 연결이 매우 자주 열리는 경우(예: 초당 여러 번) 성능에 눈에 띄는 영향을 미칠 수 있어요. 보안 연결은 HiveConf 객체에 'hive.metastore.kerberos.principal'이 올바르게 설정된 것에 의존해요.
hive-site.xml이나 사용자 지정 HiveConf에 설정된 값과 관계없이, API는 올바른 스트리밍 동작을 보장하기 위해 내부적으로 일부 설정을 재정의해요. 재정의되는 설정 목록은 다음과 같아요:
- hive.txn.manager = org.apache.hadoop.hive.ql.lockmgr.DbTxnManager
- hive.support.concurrency = true
- hive.metastore.execute.setugi = true
- hive.execution.engine = mr
I/O - 데이터 쓰기
이 클래스와 인터페이스는 트랜잭션 내에서 Hive에 데이터를 쓰는 것을 지원해요.
RecordWriter
RecordWriter는 모든 Writer가 구현하는 기본 인터페이스예요. Writer는 알려진 형식(예: CSV)의 데이터를 담은 byte[] 형태의 레코드를 받아 Hive 스트리밍이 지원하는 형식으로 쓰는 역할을 해요. RecordWriter는 필요하다면 들어오는 레코드에서 필드를 재정렬하거나 버려 Hive 테이블의 해당 컬럼에 매핑할 수 있어요. 스트리밍 클라이언트는 적절한 RecordWriter 유형을 인스턴스화해 TransactionBatch에 전달해요. 이후 스트리밍 클라이언트는 RecordWriter와 직접 상호작용하지 않아요. TransactionBatch가 이후 RecordWriter 인스턴스를 사용하고 관리하며 I/O를 수행해요. 자세한 내용은 Javadoc을 참고해요.
RecordWriter의 주요 기능은 다음과 같아요:
- 입력 레코드 수정: 입력 데이터에서 해당 테이블 컬럼이 없는 경우 필드 버리기, 특정 컬럼에 필드가 없으면 null 추가, 테이블의 필드 순서와 일치하도록 들어오는 필드 순서 변경. 이 작업은 입력 데이터 형식에 대한 이해를 요구해요. 모든 형식이(예: 데이터에 필드 이름을 포함하는 JSON은) 이 단계가 필요하진 않아요.
- 수정된 레코드 인코딩: 인코딩은 적절한 Hive SerDe를 사용한 직렬화를 포함해요.
- 레코드가 속하는 버킷 식별
- 적절한 버킷에 AcidOutputFormat의 record updater를 사용해 인코딩된 레코드를 Hive에 쓰기
DelimitedInputWriter
DelimitedInputWriter 클래스는 RecordWriter 인터페이스를 구현해요. 구분(delimited) 형식(예: CSV)의 입력 레코드를 받아 Hive에 써요. 필요하면 필드를 재정렬하고 LazySimpleSerde를 사용해 레코드를 Object로 변환한 뒤, 적절한 버킷의 기본 AcidOutputFormat의 record updater로 전달해요. Javadoc 참고.
StrictJsonWriter
StrictJsonWriter 클래스는 RecordWriter 인터페이스를 구현해요. 엄격한(strict) JSON 형식의 입력 레코드를 받아 Hive에 써요. JsonSerde를 사용해 JSON 레코드를 직접 Object로 변환한 뒤, 적절한 버킷의 기본 AcidOutputFormat의 record updater로 전달해요. Javadoc 참고.
StrictRegexWriter
StrictRegexWriter 클래스는 RecordWriter 인터페이스를 구현해요. 텍스트 형식의 입력 레코드와 정규식을 받아 Hive에 써요. 적절한 정규식을 사용해 텍스트 레코드를 RegexSerDe로 직접 Object로 변환한 뒤, 적절한 버킷의 기본 AcidOutputFormat의 record updater로 전달해요. Javadoc 참고. Hive 1.2.2+ 및 2.3.0+에서 사용 가능.
AbstractRecordWriter
스키마 조회와 레코드가 속해야 하는 버킷 계산 같은 RecordWriter 객체에 필요한 공통 코드 중 일부를 포함하는 기본 클래스예요.
오류 처리 (Error Handling)
시스템의 올바른 동작을 위해 이 API의 클라이언트가 오류를 올바르게 처리하는 것이 필수적이에요. TransactionBatch를 얻은 후 TransactionBatch에서 예외가 발생하면(SerializationError 제외) 클라이언트는 TransactionBatch.abort()를 호출해 현재 트랜잭션을 중단한 다음 TransactionBatch.close()를 호출하고 새 배치를 시작해 더 많은 데이터를 쓰거나 실패가 발생한 마지막 트랜잭션의 작업을 다시 해야 해요. 이를 따르지 않으면 드물게 파일 손상이 발생할 수 있어요. 또한 StreamingException은 이상적으로 클라이언트가 새 배치를 시작하기 전에 지수 백오프(back off)를 수행하게 해야 해요. 이 실패의 가장 가능성 높은 원인은 HDFS 과부하이므로 클러스터 안정화에 도움이 돼요.
SerializationError는 주어진 튜플을 파싱할 수 없음을 나타내요. 클라이언트는 그러한 튜플을 버리거나 데드 레터 큐(dead letter queue)로 보내도록 선택할 수 있어요. 이 예외를 본 후에도 현재 트랜잭션과 같은 TransactionBatch의 이후 트랜잭션에 더 많은 데이터를 쓸 수 있어요.
예제 - 비보안 모드 (Non-secure Mode)
///// 두 트랜잭션에서 5개 레코드 스트리밍 /////
// 가정된 HIVE 테이블 스키마:
create table alerts ( id int , msg string )
partitioned by (continent string, country string)
clustered by (id) into 5 buckets
stored as orc tblproperties("transactional"="true"); // 현재 스트리밍에는 ORC가 필요
//------- MAIN THREAD ------- //
String dbName = "testing";
String tblName = "alerts";
ArrayList<String> partitionVals = new ArrayList<String>(2);
partitionVals.add("Asia");
partitionVals.add("India");
String serdeClass = "org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe";
HiveEndPoint hiveEP = new HiveEndPoint("thrift://x.y.com:9083", dbName, tblName, partitionVals);
.. 스레드 생성 ..
//------- Thread 1 -------//
StreamingConnection connection = hiveEP.newConnection(true);
DelimitedInputWriter writer =
new DelimitedInputWriter(fieldNames,",", hiveEP);
TransactionBatch txnBatch = connection.fetchTransactionBatch(10, writer);
///// Batch 1 - 첫 번째 TXN
txnBatch.beginNextTransaction();
txnBatch.write("1,Hello streaming".getBytes());
txnBatch.write("2,Welcome to streaming".getBytes());
txnBatch.commit();
if(txnBatch.remainingTransactions() > 0) {
///// Batch 1 - 두 번째 TXN
txnBatch.beginNextTransaction();
txnBatch.write("3,Roshan Naik".getBytes());
txnBatch.write("4,Alan Gates".getBytes());
txnBatch.write("5,Owen O'Malley".getBytes());
txnBatch.commit();
txnBatch.close();
connection.close();
}
txnBatch = connection.fetchTransactionBatch(10, writer);
///// Batch 2 - 첫 번째 TXN
txnBatch.beginNextTransaction();
txnBatch.write("6,David Schorow".getBytes());
txnBatch.write("7,Sushant Sowmyan".getBytes());
txnBatch.commit();
if(txnBatch.remainingTransactions() > 0) {
///// Batch 2 - 두 번째 TXN
txnBatch.beginNextTransaction();
txnBatch.write("8,Ashutosh Chauhan".getBytes());
txnBatch.write("9,Thejas Nair".getBytes());
txnBatch.commit();
txnBatch.close();
}
connection.close();
//------- Thread 2 -------//
StreamingConnection connection2 = hiveEP.newConnection(true);
DelimitedInputWriter writer2 =
new DelimitedInputWriter(fieldNames,",", hiveEP);
TransactionBatch txnBatch2= connection.fetchTransactionBatch(10, writer2);
///// Batch 1 - 첫 번째 TXN
txnBatch2.beginNextTransaction();
txnBatch2.write("21,Venkat Ranganathan".getBytes());
txnBatch2.write("22,Bowen Zhang".getBytes());
txnBatch2.commit();
///// Batch 1 - 두 번째 TXN
txnBatch2.beginNextTransaction();
txnBatch2.write("23,Venkatesh Seetaram".getBytes());
txnBatch2.write("24,Deepesh Khandelwal".getBytes());
txnBatch2.commit();
txnBatch2.close();
connection.close();
txnBatch = connection.fetchTransactionBatch(10, writer);
///// Batch 2 - 첫 번째 TXN
txnBatch.beginNextTransaction();
txnBatch.write("26,David Schorow".getBytes());
txnBatch.write("27,Sushant Sowmyan".getBytes());
txnBatch.commit();
txnBatch2.close();
connection2.close();
예제 - 보안 스트리밍 (Secure Streaming)
Kerberos로 보안 Hive metastore에 연결하려면 UserGroupInformation(UGI) 객체가 필요해요. 이 UGI 객체는 외부에서 획득해 EndPoint.newConnection의 인자로 전달해야 해요. 그 연결 객체를 사용해 수행되는 모든 후속 내부 연산(트랜잭션 배치 획득, 쓰기, 커밋 등)은 필요에 따라 내부적으로 ugi.doAs 블록으로 자동 래핑돼요.
중요: Kerberos로 연결하려면 EndPoint.newConnection()의 'authenticatedUser' 인자가 Kerberos 로그인을 수행하는 데 사용되었어야 해요. 또한 'hive.metastore.kerberos.principal' 설정이 hive-site.xml 또는 'conf' 인자(null이 아닌 경우)에 올바르게 설정되어야 해요. hive-site.xml을 사용하는 경우 그 디렉토리가 클래스패스에 포함되어야 해요.
import org.apache.hadoop.security.UserGroupInformation;
HiveEndPoint hiveEP2 = ... ;
UserGroupInformation ugi = .. authenticateWithKerberos(principal,keytab);
StreamingConnection secureConn = hiveEP2.newConnection(true, null, ugi);
DelimitedInputWriter writer3 = new DelimitedInputWriter(fieldNames, ",", hiveEP2);
TransactionBatch txnBatch3= secureConn.fetchTransactionBatch(10, writer3);
///// Batch 1 - 첫 번째 TXN - 보안 연결 통해
txnBatch3.beginNextTransaction();
txnBatch3.write("28,Eric Baldeschwieler".getBytes());
txnBatch3.write("29,Ari Zilka".getBytes());
txnBatch3.commit();
txnBatch3.close();
secureConn.close();
지식 베이스 (Knowledge Base)
더 알아보기 (Learn more)
- Hive Transactions에서 ACID와 컴팩션에 대해 확인할 수 있어요.
- HCatalog Streaming Mutation API 및 Hive 3 Streaming API 문서를 볼 수 있어요.