Apache Storm용 Pulsar 어댑터
Apache Storm용 Pulsar 어댑터
Pulsar Storm은 Apache Storm 토폴로지(topology)와 통합하기 위한 어댑터예요. 데이터를 보내고 받기 위한 핵심 Storm 구현을 제공해요. 애플리케이션은 범용 Pulsar 스파우트(spout)를 통해 Storm 토폴로지에 데이터를 주입하고, 범용 Pulsar 볼트(bolt)를 통해 Storm 토폴로지의 데이터를 소비할 수 있어요.
출처: 문서
본문
Pulsar Storm 어댑터 사용하기 (Using the Pulsar Storm Adaptor)
Pulsar Storm 어댑터를 사용하려면 Pulsar Storm 어댑터 의존성을 포함해야 해요.
<dependency>
<groupId>org.apache.pulsar</groupId>
<artifactId>pulsar-storm</artifactId>
<version>${pulsar.version}</version>
</dependency>
Pulsar 스파우트 (Pulsar Spout)
Pulsar 스파우트는 토픽에 게시된 데이터를 Storm 토폴로지가 소비할 수 있게 해줘요. 수신한 메시지와 클라이언트가 제공한 MessageToValuesMapper에 기반해 Storm 튜플(tuple)을 방출(emit)해요.
다운스트림 볼트가 처리하지 못한 튜플은 스파우트가 지수 백오프(exponential backoff)로 재주입해요. 설정 가능한 타임아웃(기본 60초) 또는 설정 가능한 재시도 횟수 중 먼저 도달하는 것까지 재시도하고, 그 후 컨슈머가 ack 처리해요. 스파우트 구성 예제는 다음과 같아요.
MessageToValuesMapper messageToValuesMapper = new MessageToValuesMapper() {
@Override
public Values toValues(Message msg) {
return new Values(new String(msg.getData()));
}
@Override
public void declareOutputFields(OutputFieldsDeclarer declarer) {
// declare the output fields
declarer.declare(new Fields("string"));
}
};
// Configure a Pulsar Spout
PulsarSpoutConfiguration spoutConf = new PulsarSpoutConfiguration();
spoutConf.setServiceUrl("pulsar://broker.messaging.usw.example.com:6650");
spoutConf.setTopic("persistent://my-property/usw/my-ns/my-topic1");
spoutConf.setSubscriptionName("my-subscriber-name1");
spoutConf.setMessageToValuesMapper(messageToValuesMapper);
// Create a Pulsar Spout
PulsarSpout spout = new PulsarSpout(spoutConf);
전체 예제는 여기를 클릭해서 볼 수 있어요.
Pulsar 볼트 (Pulsar Bolt)
Pulsar 볼트는 Storm 토폴로지의 데이터를 토픽에 게시할 수 있게 해줘요. 수신한 Storm 튜플과 클라이언트가 제공한 TupleToMessageMapper에 기반해 메시지를 게시해요.
파티션 토픽을 사용해 서로 다른 토픽에 메시지를 게시할 수도 있어요. TupleToMessageMapper 구현에서 메시지에 "key"를 제공하면 같은 키를 가진 메시지가 같은 토픽으로 전송돼요. 볼트 예제는 다음과 같아요.
TupleToMessageMapper tupleToMessageMapper = new TupleToMessageMapper() {
@Override
public TypedMessageBuilder<byte[]> toMessage(TypedMessageBuilder<byte[]> msgBuilder, Tuple tuple) {
String receivedMessage = tuple.getString(0);
// message processing
String processedMsg = receivedMessage + "-processed";
return msgBuilder.value(processedMsg.getBytes());
}
@Override
public void declareOutputFields(OutputFieldsDeclarer declarer) {
// declare the output fields
}
};
// Configure a Pulsar Bolt
PulsarBoltConfiguration boltConf = new PulsarBoltConfiguration();
boltConf.setServiceUrl("pulsar://broker.messaging.usw.example.com:6650");
boltConf.setTopic("persistent://my-property/usw/my-ns/my-topic2");
boltConf.setTupleToMessageMapper(tupleToMessageMapper);
// Create a Pulsar Bolt
PulsarBolt bolt = new PulsarBolt(boltConf);
더 알아보기 (Learn more)
- Apache Storm의 기본 개념은 Storm 공식 문서를 참고해요.
- Pulsar 토픽과 구독에 대해 더 알고 싶다면 핵심 개념 문서를 살펴보세요.
- 파티션 토픽 동작 방식이 궁금하다면 파티션 토픽 관련 문서를 참고해요.