HCatalog 알림
HCatalog 알림 (Notification)
HCatalog가 일부 시스템 이벤트에 대해 알림을 보내는 기능을 설명해요. 새 파티션이 추가될 때 메시지 버스(JMS)를 통해 통지받아 Oozie 같은 애플리케이션이 후속 작업을 예약할 수 있게 합니다.
출처: 문서
본문
개요 (Overview)
HCatalog는 0.2 버전부터 시스템에서 발생하는 특정 이벤트에 대한 알림을 제공합니다. 이렇게 하여 Oozie 같은 애플리케이션은 그 이벤트를 기다렸다가 그에 의존하는 작업을 예약할 수 있습니다. 현재 HCatalog 버전은 두 종류의 이벤트를 지원합니다:
- 새 파티션이 추가될 때 알림
- 파티션 집합이 추가될 때 알림
새 파티션이 추가될 때 알림을 보내는 데는 추가 작업이 필요 없습니다. 기존 addPartition 호출이 알림 메시지를 보냅니다.
새 파티션에 대한 알림 (Notification for a New Partition)
새 파티션이 추가되었다는 알림을 받으려면 다음 세 단계를 따라야 합니다.
- 메시지 수신을 시작하려면 여기에 나온 것처럼 메시지 버스에 대한 연결을 만듭니다:
ConnectionFactory connFac = new ActiveMQConnectionFactory(amqurl);
Connection conn = connFac.createConnection();
conn.start();
- 관심 있는 토픽을 구독합니다. 메시지 버스에서 구독할 때, 해당 토픽으로 전달되는 메시지를 받으려면 특정 토픽에 구독해야 합니다.
특정 테이블에 해당하는 토픽 이름은 테이블 속성에 저장되어 있으며 다음 코드 조각으로 검색할 수 있습니다:
HiveMetaStoreClient msc = new HiveMetaStoreClient(hiveConf);
String topicName = msc.getTable("mydb",
"myTbl").getParameters().get(HCatConstants.HCAT_MSGBUS_TOPIC_NAME);
- 다음과 같이 토픽 이름으로 토픽을 구독합니다:
Session session = conn.createSession(true, Session.SESSION_TRANSACTED);
Destination hcatTopic = session.createTopic(topicName);
MessageConsumer consumer = session.createConsumer(hcatTopic);
consumer.setMessageListener(this);
- 메시지 수신을 시작하려면 JMS 인터페이스
MessageListener를 구현해야 하며, 그러면onMessage(Message msg)메서드를 구현하게 됩니다. 이 메서드는 메시지 버스에 새 메시지가 도착할 때마다 호출됩니다. 메시지는 해당 파티션을 나타내는 파티션 객체를 포함하며, 아래처럼 검색할 수 있습니다:
@Override
public void onMessage(Message msg) {
// 이 테이블에서는 add_partition 이벤트에만 관심이 있다.
// 그래서 먼저 메시지 타입을 확인한다.
if(msg.getStringProperty(HCatConstants.HCAT_EVENT).equals(HCatConstants.HCAT_ADD_PARTITION_EVENT)){
Object obj = (((ObjectMessage)msg).getObject());
}
}
이것이 동작하려면 클래스패스에 JMS jar가 있어야 합니다. 추가로 JMS 프로바이더의 jar도 클래스패스에 있어야 합니다. HCatalog는 JMS 프로바이더로 ActiveMQ와 함께 테스트되었지만, 어떤 JMS 프로바이더든 사용할 수 있습니다. ActiveMQ는 http://activemq.apache.org/activemq-550-release.html에서 얻을 수 있습니다.
파티션 집합에 대한 알림 (Notification for a Set of Partitions)
때로는 다른 연산을 진행하기 전에 파티션 모음이 끝나기를 기다려야 합니다. 예를 들어 하루에 대한 모든 파티션이 끝난 뒤 처리를 시작하고 싶을 수 있습니다. 그러나 HCatalog에는 파티션의 모음이나 계층(hierarchy) 개념이 없습니다. 이를 지원하기 위해 HCatalog는 데이터 작성자(writer)가 파티션 모음 작성을 끝냈을 때 신호를 보낼 수 있게 합니다. 데이터 읽는 사람(reader)은 읽기를 시작하기 전에 이 신호를 기다릴 수 있습니다.
아래 예시 코드는 파티션 집합이 추가되었을 때 알림을 보내는 방법을 보여줍니다.
신호를 보내기 위해 데이터 작성자는 이렇게 합니다:
HiveMetaStoreClient msc = new HiveMetaStoreClient(conf);
// 파티션 키 이름과 값을 지정하는 맵 생성
Map<String,String> partMap = new HashMap<String, String>();
partMap.put("date","20110711");
partMap.put("country","*");
// 파티션을 "완료"로 표시
msc.markPartitionForEvent("mydb", "mytbl", partMap, PartitionEventType.LOAD_DONE);
이 알림을 받으려면 소비자는 다음을 수행해야 합니다:
- 위의 1단계와 2단계를 반복해 알림 시스템에 연결하고 토픽을 구독합니다.
- 다음 예시처럼 알림을 받습니다:
HiveMetaStoreClient msc = new HiveMetaStoreClient(conf);
// 파티션 키 이름과 값을 지정하는 맵 생성
Map<String,String> partMap = new HashMap<String, String>();
partMap.put("date","20110711");
partMap.put("country","*");
// 파티션을 "완료"로 표시
msc.markPartitionForEvent("mydb", "mytbl", partMap, PartitionEventType.LOAD_DONE);
소비자가 메시지 버스에 등록되어 현재 활성 상태라면, 프로듀서가 파티션을 "완료"로 표시하면 메시지 버스로부터 콜백을 받게 됩니다. 또는 소비자는 특정 파티션을 메타스토어에 명시적으로 조회할 수도 있습니다. 다음 코드는 소비자 관점의 사용을 보여줍니다:
// 특정 파티션이 표시되었는지 메타스토어에 조회한다.
boolean marked = msc.isPartitionMarkedForEvent("mydb", "mytbl", partMap, PartitionEventType.LOAD_DONE);
// 또는 메시지 버스에 등록하고 비동기 콜백을 받는다.
ConnectionFactory connFac = new ActiveMQConnectionFactory(amqurl);
Connection conn = connFac.createConnection();
conn.start();
Session session = conn.createSession(true, Session.SESSION_TRANSACTED);
Destination hcatTopic = session.createTopic(topic);
MessageConsumer consumer = session.createConsumer(hcatTopic);
consumer.setMessageListener(this);
public void onMessage(Message msg) {
MapMessage mapMsg = (MapMessage)msg;
Enumeration<String> keys = mapMsg.getMapNames();
// 모든 키를 열거한다. 완료로 표시된 특정 파티션을
// 지정하는 키-값 쌍을 출력한다. 이 경우 출력은:
// date : 20110711
// country: *
while(keys.hasMoreElements()){
String key = keys.nextElement();
System.out.println(key + " : " + mapMsg.getString(key));
}
System.out.println("Message: "+msg);
서버 구성 (Server Configuration)
알림을 활성화하려면 서버를 구성해야 합니다(아래 참고).
알림을 비활성화하려면 hive.metastore.event.listeners를 비워 두거나 hive-site.xml에서 제거하면 됩니다.
JMS 알림 활성화 (Enable JMS Notifications)
HCatalog 서버의 hive-site.xml 파일에 다음 변경 사항을 만들어(추가/수정) 알림을 켭니다.
<property>
<name>hive.metastore.event.expiry.duration</name>
<value>300L</value>
<description>Duration after which events expire from events table (in seconds)</description>
</property>
<property>
<name>hive.metastore.event.clean.freq</name>
<value>360L</value>
<description>Frequency at which timer task runs to purge expired events in metastore (in seconds).</description>
</property>
<property>
<name>msgbus.brokerurl</name>
<value>tcp://localhost:61616</value>
<description></description>
</property>
<property>
<name>msgbus.username</name>
<value></value>
<description></description>
</property>
<property>
<name>msgbus.password</name>
<value></value>
<description></description>
</property>
서버가 알림 지원으로 시작하려면 다음이 클래스패스에 있어야 합니다:
(a) activemq jar
(b) 알림에 적합하게 구성된 속성이 있는 jndi.properties 파일
그런 다음 환경을 설정하려면 다음 지침을 따르세요:
- HCatalog 서버 시작 스크립트는 $YOUR_HCAT_SERVER
/share/hcatalog/scripts/hcat_server_start.sh입니다. - 이 스크립트는 AUX_CLASSPATH 환경 변수로 클래스패스가 설정되기를 기대합니다.
- 따라서 (a)와 (b)를 충족하도록 AUX_CLASSPATH를 설정하세요.
jndi.properties파일은 $YOUR_HCAT_SERVER/etc/hcatalog/jndi.properties에 있습니다.jndi.properties파일에서 다음 속성을 주석 해제하고 설정해야 합니다:
java.naming.factory.initial = org.apache.activemq.jndi.ActiveMQInitialContextFactory
java.naming.provider.url = tcp://localhost:61616(이것은 설정의 ActiveMQ URL입니다.)
토픽 이름 (Topic Names)
서버가 알림용으로 구성된 동안 테이블이 생성되면, 기본 토픽 이름이 테이블 속성으로 자동 설정됩니다. 이전에 생성된 테이블(다른 HCatalog 설치 또는 현재 설치에서 알림을 활성화하기 전에 생성된)에 알림을 사용하려면 토픽 이름을 수동으로 설정해야 합니다. 예를 들어:
$YOUR_HCAT_CLIENT_HOME/bin/hcat -e "ALTER TABLE access SET hcat.msgbus.topic.name=$TOPIC_NAME"
그런 다음 ActiveMQ 소비자(들)를 $TOPIC_NAME에서 준 토픽의 메시지를 듣도록 구성해야 합니다. 좋은 기본 정책은 TOPIC_NAME = "$database.$table"(말 그대로 점)입니다.
더 알아보기 (Learn more)
HCatalog 알림은 새 파티션 추가 이벤트를 JMS 메시지 버스(ActiveMQ 등)로 전파해요. addPartition 호출에 자동으로 알림이 발생하고, markPartitionForEvent로 파티션 집합 완료를 표시할 수 있습니다. hive-site.xml의 metastore 이벤트 속성과 jndi.properties로 서버를 구성합니다.