스트리밍 데이터 수집 V2

스트리밍 데이터 수집 V2 (Streaming Data Ingest V2)

Hive 3.0.0부터 Streaming Data Ingest는 deprecated 되고 새 V2 API(HIVE-19205)로 대체됐어요. 이 API로 NiFi, Flume 같은 스트리밍 클라이언트가 Hive에 데이터를 계속 흘려보낼 수 있어요.

출처: 문서

본문

Hive 3.0.0 릴리스부터 Streaming Data Ingest는 deprecated 되고 더 새로운 V2 API(HIVE-19205)로 대체됩니다.

Hive 스트리밍 API

전통적으로 Hive에 새 데이터를 추가하려면 많은 양의 데이터를 HDFS에 모은 다음 주기적으로 새 파티션을 추가해야 합니다. 이는 본질적으로 "배치 삽입(batch insertion)"입니다.

Hive 스트리밍 API는 데이터를 Hive에 지속적으로 흘려보낼 수 있게 해줍니다. 들어오는 데이터는 기존 Hive 파티션이나 테이블에 소규모 레코드 배치로 지속적으로 커밋될 수 있습니다. 데이터가 커밋되면 이후에 시작된 모든 Hive 쿼리에 즉시 표시됩니다.

이 API는 데이터를 지속적으로 생성하는 NiFi, Flume, Storm 같은 스트리밍 클라이언트를 위한 것입니다. 스트리밍 지원은 Hive의 ACID 기반 insert/update 지원 위에 구축됩니다(Hive Transactions 참고).

Hive 스트리밍 API의 클래스와 인터페이스는 크게 두 세트로 분류됩니다. 첫 번째 세트는 연결과 트랜잭션 관리를 지원하고, 두 번째 세트는 I/O를 지원합니다. 트랜잭션은 metastore가 관리합니다. 쓰기는 테이블이 정의한 대상 파일시스템(HDFS, S3A 등)에 직접 수행됩니다.

비파티션 테이블, 정적 파티션이 있는 파티션 테이블, 동적 파티션이 있는 파티션 테이블로의 스트리밍이 모두 지원됩니다. API는 Kerberos 인증과 Storage 기반 권한 부여를 지원합니다. 클라이언트 사용자는 API 호출 전에 kerberos로 로그인해야 합니다. 로그인한 사용자는 대상 파티션이나 테이블 위치에 쓰기 위한 적절한 저장소 권한이 있어야 합니다. 'hive' 사용자를 사용하는 것이 권장됩니다. 그렇게 하면 hive 쿼리가 doAs를 false로 설정(쿼리가 hive 사용자로 실행됨)한 상태로 (스트리밍 API가 쓴) 데이터를 다시 읽을 수 있기 때문입니다.

패키징 참고: API는 Java 패키지 org.apache.hive.streaming에 정의되어 있으며 Hive의 hive-streaming Maven 모듈의 일부입니다.

Streaming Mutation API의 Deprecation 및 제거

3.0.0 릴리스부터 Hive는 hive-hcatalog-streaming 모듈의 Streaming Mutation API를 deprecated 처리했으며 향후 릴리스에서 더 이상 지원하지 않습니다. 새 hive-streaming 모듈은 더 이상 mutation API를 지원하지 않습니다.

스트리밍 요구 사항

스트리밍을 사용하려면 몇 가지가 필요합니다.

  1. hive-site.xml에서 스트리밍용 ACID 지원을 활성화하려면 다음 설정이 필요합니다:
    1. hive.txn.manager = org.apache.hadoop.hive.ql.lockmgr.DbTxnManager
    2. hive.compactor.initiator.on = true (더 중요한 세부사항은 여기 참고)
    3. hive.compactor.cleaner.on = true (Hive 4.0.0 이후부터. 자세한 내용은 여기 참고)
    4. hive.compactor.worker.threads > 0
  2. 테이블 생성 시 ""stored as orc"" 를 지정해야 합니다. 현재 ORC 저장 형식만 지원됩니다.
  3. 생성 시 테이블에 tblproperties("transactional"="true")를 설정해야 합니다.
  4. 클라이언트 스트리밍 프로세스의 사용자는 테이블·파티션에 쓰고 테이블에 파티션을 만들 수 있는 필요한 권한이 있어야 합니다.

제한 사항

기본적으로 현재 스트리밍 API는 구분자(delimited) 입력 데이터(CSV, 탭 구분 등), JSON, Regex 형식 데이터의 스트리밍만 지원합니다. Hive 3.0.0 릴리스에서 지원되는 레코드 작성기는 다음과 같습니다:

StrictDelimitedInputWriter

StrictJsonWriter

StrictRegexWriter

이 모든 레코드 작성기는 엄격한 스키마 일치를 기대합니다. 즉 레코드의 스키마가 테이블 스키마와 정확히 일치해야 합니다(참고: 작성기는 스키마 검사를 수행하지 않으며, 레코드 스키마가 테이블 스키마와 일치하는지 확인하는 것은 클라이언트의 몫입니다).

다른 입력 형식에 대한 지원은 RecordWriter 인터페이스의 추가 구현으로 제공할 수 있습니다.

현재 대상 테이블의 형식으로는 ORC만 지원됩니다.

API 사용법

트랜잭션 및 연결 관리

HiveStreamingConnection

HiveStreamingConnection 클래스는 스트리밍 연결 관련 정보를 설명합니다. 데이터베이스, 테이블, 파티션 이름, 연결할 metastore URI, 대상 파티션이나 테이블로 레코드를 스트리밍하는 데 사용할 레코드 작성기를 설명합니다. 이 모든 정보는 Builder API로 지정할 수 있으며, 그러면 스트리밍 목적으로 Hive MetaStore에 연결을 수립합니다. Builder API에서 connect를 호출하면 StreamingConnection 객체를 반환합니다. 그런 다음 StreamingConnection을 사용해 I/O를 수행하기 위한 새 트랜잭션을 시작할 수 있습니다.

  • 동시성 주의: 스트리밍 연결 API와 레코드 작성기 API는 스레드 안전하지 않습니다. 스트리밍 연결 생성, 트랜잭션 begin/commit/abort, write와 close는 같은 스레드에서 호출되어야 합니다. close() 또는 abortTransaction()을 별도 스레드에서 트리거해야 한다면 외부 변수나 동기화 메커니즘을 통해 조정해야 합니다.

HiveStreamingConnection API는 또한 2가지 파티셔닝 모드(정적 vs 동적)를 지원합니다. 정적 파티셔닝 모드에서는 builder API를 통해 파티션 컬럼 값을 미리 지정할 수 있습니다. 정적 파티션이 존재하면(사용자가 미리 만들었거나 테이블에 이미 존재하는 파티션) 스트리밍 연결이 그것을 사용하고, 존재하지 않으면 지정된 값을 사용해 새 정적 파티션을 만듭니다. 테이블이 파티션되어 있고 builder API에 정적 파티션 값을 지정하지 않으면 hive 스트리밍 연결은 동적 파티셔닝 모드를 사용하며, 이 모드에서 hive 스트리밍 연결은 파티션 값이 레코드의 마지막 컬럼이기를 기대합니다(hive 동적 파티셔닝이 동작하는 방식과 유사). 예를 들어 테이블이 2개 컬럼(year와 month)으로 파티션되면 hive 스트리밍 연결은 입력 레코드에서 마지막 2개 컬럼을 추출하고, 이를 사용해 metastore에 파티션이 동적으로 생성됩니다.

트랜잭션은 전통적인 데이터베이스 시스템과 약간 다르게 구현됩니다. 각 트랜잭션에는 id가 있고 여러 트랜잭션이 "Transaction Batch"로 그룹화됩니다. 이는 여러 트랜잭션의 레코드를 더 적은 파일로 그룹화하는 데 도움이 됩니다(트랜잭션당 1개 파일 대신). hive 스트리밍 연결 생성 중에 builder API로 트랜잭션 배치 크기를 지정할 수 있습니다. 트랜잭션 관리는 API 뒤에 완전히 숨겨져 있어 대부분의 경우 사용자가 트랜잭션 배치 크기 조정을 걱정할 필요가 없습니다(이것은 전문가 수준 설정이며 향후 릴리스에서 존중되지 않을 수 있습니다). 또한 현재 트랜잭션 배치가 소진되면 API는 beginTransaction() 호출 시 자동으로 다음 트랜잭션 배치로 넘어갑니다. 트랜잭션 배치 크기를 기본값 1로 두고 각 트랜잭션 아래에 수천 개의 레코드를 함께 그룹화하는 것이 권장됩니다. 각 트랜잭션은 파일시스템의 delta 디렉터리에 해당하므로, 트랜잭션을 너무 자주 커밋하면 너무 많은 작은 디렉터리가 생성될 수 있습니다.

TransactionBatch의 트랜잭션은 hive.txn.timeout 초 후에 커밋되거나 중단되지 않으면 결국 Metastore에 의해 만료됩니다. 트랜잭션을 활성 상태로 유지하기 위해 HiveStreamingConnection에는 heartbeat 스레드가 있으며, 기본적으로 모든 열린 트랜잭션에 대해 (hive.txn.timeout/2) 간격으로 heartbeat를 보냅니다.

자세한 내용은 HiveStreamingConnection Javadoc을 참고하세요.

사용 지침

일반적으로 각 트랜잭션에 더 많은 레코드가 포함될수록 더 많은 처리량을 달성할 수 있습니다. 특정 레코드 수 후 또는 특정 시간 간격 후 중 먼저 도래하는 시점에 커밋하는 것이 일반적입니다. 후자는 이벤트 흐름 속도가 변할 때 트랜잭션이 너무 오래 열려 있지 않도록 보장합니다. 단일 트랜잭션에 포함할 수 있는 데이터 양에는 실질적인 제한이 없습니다. 유일한 우려는 트랜잭션이 실패할 경우 다시 재생해야 할 데이터 양입니다. TransactionBatch의 개념은 HiveStreamingConnection API가 파일시스템에 만드는 파일(및 delta 디렉터리)의 수를 줄이는 역할을 합니다. 주어진 트랜잭션 배치의 모든 트랜잭션은 (버킷별로) 같은 물리 파일에 쓰므로, 파티션은 열린 트랜잭션을 포함한 어떤 배치의 가장 이른 트랜잭션 수준까지만 컴팩트될 수 있습니다. 따라서 TransactionBatch를 지나치게 크게 만들지 않아야 합니다. 일정 시간 후에(사용되지 않은 트랜잭션이 있더라도) TransactionBatch를 닫는 타이머를 포함하는 것이 합리적입니다.

HiveStreamingConnection은 쓰기 처리량에 매우 최적화되어 있으며(Delta Streaming Optimizations) 그 결과 Hive 스트리밍 수집이 생성하는 delta 파일은 높은 처리량 쓰기를 용이하게 하기 위해 많은 ORC 기능(딕셔너리 인코딩, 인덱스, 압축 등)이 비활성화되어 있습니다. 컴팩터가 작동하면 이 delta 파일들은 읽기·저장 최적화 ORC 형식(딕셔너리 인코딩, 인덱스, 압축 활성화)으로 다시 쓰여집니다. 따라서 컴팩터를 더 공격적으로/자주 구성하는 것이 권장됩니다(Compactor 참고). 그래야 컴팩트되고 최적화된 ORC 파일이 생성됩니다.

HiveConf 객체에 대한 참고

HiveStreamingConnect builder API는 HiveConf 인자를 받습니다. 이것은 null로 설정하거나 미리 만든 HiveConf 객체를 제공할 수 있습니다. null이면 연결을 위해 HiveConf 객체가 내부적으로 생성·사용됩니다. HiveConf 객체가 인스턴스화될 때, hive-site.xml을 포함하는 디렉터리가 java classpath의 일부이면 HiveConf 객체는 그 값으로 초기화됩니다. hive-site.xml이 없으면 기본값으로 초기화됩니다. 객체를 미리 만들고 여러 연결에 걸쳐 재사용하면, 연결을 매우 자주 열 때(예: 초당 여러 번) 성능에 눈에 띄는 영향을 줄 수 있습니다. 보안 연결은 HiveConf 객체에 '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.exec.dynamic.partition.mode = nonstrict
  • hive.exec.orc.delta.streaming.optimizations.enabled = true
  • hive.metastore.client.cache.enabled = false

I/O – 데이터 쓰기

이 클래스와 인터페이스는 트랜잭션 내에서 Hive에 데이터를 쓰는 지원을 제공합니다.

RecordWriter

RecordWriter는 모든 Writer가 구현하는 기본 인터페이스입니다. Writer는 알려진 형식(예: CSV)의 데이터를 포함하는 byte[] (또는 구성 가능한 줄 구분자가 있는 InputStream) 형태의 레코드를 받아 Hive 스트리밍이 지원하는 형식으로 쓰는 역할을 합니다. Strict 구현의 RecordWriter는 레코드 스키마가 테이블 스키마와 정확히 일치하기를 기대합니다. 동적 파티셔닝 모드로 쓰는 RecordWriter는 파티션 컬럼이 각 레코드의 마지막 컬럼이기를 기대합니다. 파티션 컬럼 값이 비어 있거나 null이면 레코드는 HIVE_DEFAULT_PARTITION 으로 갑니다. 스트리밍 클라이언트는 적절한 RecordWriter 타입을 인스턴스화해 HiveStreamingConnection builder API에 전달합니다. 이후 스트리밍 클라이언트는 RecordWriter와 직접 상호작용하지 않습니다. 그 후 StreamingConnection 객체가 RecordWriter 인스턴스를 사용·관리해 I/O를 수행합니다. 자세한 내용은 Javadoc을 참고하세요.

RecordWriter의 주요 기능은 다음과 같습니다:

  1. 입력 레코드 수정: 해당 테이블 컬럼이 없으면 입력 데이터에서 필드를 드롭하고, 특정 컬럼에 필드가 없으면 null을 추가하고, 파티션 컬럼 값이 null이거나 비어 있으면 HIVE_DEFAULT_PARTITION을 추가하는 작업이 포함될 수 있습니다. 동적으로 파티션을 만드는 것은 파티션 값을 추출하기 위해 마지막 컬럼을 뽑아내려면 들어오는 데이터 형식을 이해해야 합니다.
  2. 수정된 레코드 인코딩: 인코딩은 적절한 Hive SerDe를 사용한 직렬화를 포함합니다.
  3. bucket 테이블의 경우, 레코드에서 bucket 컬럼 값을 추출해 레코드가 속한 버킷을 식별합니다.
  4. 파티션 테이블의 경우, 동적 파티셔닝 모드에서 레코드의 마지막 N개 컬럼(N은 파티션 수)에서 파티션 컬럼 값을 추출해 레코드가 속한 파티션을 식별합니다.
  5. 적절한 버킷에 대해 AcidOutputFormat의 record updater를 사용해 인코딩된 레코드를 Hive에 씁니다.

StrictDelimitedInputWriter

StrictDelimitedInputWriter 클래스는 RecordWriter 인터페이스를 구현합니다. 구분자 형식(예: CSV)의 입력 레코드를 받아 Hive에 씁니다. 레코드 스키마가 테이블 스키마와 일치하고 파티션 값이 마지막에 있기를 기대합니다. 입력 레코드는 LazySimpleSerde를 사용해 Object로 변환되어 bucket과 파티션 컬럼을 추출하고, 그런 다음 적절한 버킷의 기본 AcidOutputFormat record updater에 전달됩니다. Javadoc 참고.

StrictJsonWriter

StrictJsonWriter 클래스는 RecordWriter 인터페이스를 구현합니다. 엄격한 JSON 형식의 입력 레코드를 받아 Hive에 씁니다. JsonSerde를 사용해 JSON 레코드를 Object로 직접 변환하고, 그런 다음 적절한 버킷과 파티션의 기본 AcidOutputFormat record updater에 전달됩니다. Javadoc 참고.

StrictRegexWriter

StrictRegexWriter 클래스는 RecordWriter 인터페이스를 구현합니다. 텍스트 형식의 입력 레코드와 regex를 받아 Hive에 씁니다. 적절한 regex를 사용해 텍스트 레코드를 RegexSerDe로 Object로 직접 변환하고, 그런 다음 적절한 버킷의 기본 AcidOutputFormat record updater에 전달됩니다. Javadoc 참고.

AbstractRecordWriter

이것은 스키마 조회와 레코드가 속해야 하는 bucket·파티션 계산 같은 RecordWriter 객체에 필요한 공통 코드의 일부를 포함하는 기본 클래스입니다.

오류 처리

시스템이 제대로 기능하려면 이 API의 클라이언트가 오류를 올바르게 처리하는 것이 필수입니다. API는 StreamingException만 반환하지만, 서로 다른 상황에서 던져지는 StreamingException의 여러 하위 클래스가 있습니다.

  • ConnectionError - metastore에 연결을 수립할 수 없거나 HiveStreamingConnection connect API를 잘못 사용했을 때
  • InvalidTable - 대상 테이블이 존재하지 않거나 ACID 트랜잭션 테이블이 아닐 때
  • InvalidTransactionState - 트랜잭션 배치가 잘못된 상태가 될 때
  • SerializationError - SerDe가 레코드의 직렬화/역직렬화/쓰기 중에 예외를 던질 때
  • StreamingIOFailure - 대상 파티션을 만들 수 없거나, record updater가 쓰기·플러시 중 IO 오류를 던질 때
  • TransactionError - 내부 트랜잭션을 커밋하거나 중단할 수 없을 때

어떤 예외를 재시도하고(일반적으로 약간의 백오프를 두고), 무시하고, 다시 던질지 결정하는 것은 클라이언트의 몫입니다. 일반적으로 연결 관련 예외는 지수 백오프로 재시도할 수 있습니다. 직렬화 관련 오류는 던지거나 무시할 수 있습니다(일부 입력 레코드가 잘못되었거나 손상된 경우 드롭할 수 있음).

예제

///// 두 트랜잭션에서 다섯 레코드 스트리밍 /////
 
// 가정한 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";

.. spin up thread 1 ..
// static partition values
ArrayList<String> partitionVals = new ArrayList<String>(2);
partitionVals.add("Asia");
partitionVals.add("India");

// create delimited record writer whose schema exactly matches table schema
StrictDelimitedInputWriter writer = StrictDelimitedInputWriter.newBuilder()
                                      .withFieldDelimiter(',')
                                      .build();
// create and open streaming connection (default.src table has to exist already)
StreamingConnection connection = HiveStreamingConnection.newBuilder()
                                    .withDatabase(dbName)
                                    .withTable(tblName)
                                    .withStaticPartitionValues(partitionVals)
                                    .withAgentInfo("example-agent-1")
                                    .withRecordWriter(writer)
                                    .withHiveConf(hiveConf)
                                    .connect();
// begin a transaction, write records and commit 1st transaction
connection.beginTransaction();
connection.write("1,val1".getBytes());
connection.write("2,val2".getBytes());
connection.commitTransaction();
// begin another transaction, write more records and commit 2nd transaction
connection.beginTransaction();
connection.write("3,val3".getBytes());
connection.write("4,val4".getBytes());
connection.commitTransaction();
// close the streaming connection
connection.close();

.. spin up thread 2 ..
// dynamic partitioning
// create delimited record writer whose schema exactly matches table schema
StrictDelimitedInputWriter writer = StrictDelimitedInputWriter.newBuilder()
                                      .withFieldDelimiter(',')
                                      .build();
// create and open streaming connection (default.src table has to exist already)
StreamingConnection connection = HiveStreamingConnection.newBuilder()
                                    .withDatabase(dbName)
                                    .withTable(tblName)
                                    .withAgentInfo("example-agent-1")
                                    .withRecordWriter(writer)
                                    .withHiveConf(hiveConf)
                                    .connect();
// begin a transaction, write records and commit 1st transaction
connection.beginTransaction();
// dynamic partition mode where last 2 columns are partition values
connection.write("11,val11,Asia,China".getBytes());
connection.write("12,val12,Asia,India".getBytes());
connection.commitTransaction();
// begin another transaction, write more records and commit 2nd transaction
connection.beginTransaction();
connection.write("13,val13,Europe,Germany".getBytes());
connection.write("14,val14,Asia,India".getBytes());
connection.commitTransaction();
// close the streaming connection
connection.close();

더 알아보기 (Learn more)

스트리밍은 ACID 트랜잭션 위에 구축되므로 Hive Transactions 문서를 먼저 이해하는 것이 좋아요. 레코드를 직렬화하는 SerDe 개념은 SerDe 문서를 참고하세요.