RabbitMQ 커넥터
RabbitMQ 커넥터
RabbitMQ에서 데이터 스트림에 접근할 수 있게 해주는 DataStream 커넥터예요.
출처: 문서
본문
RabbitMQ 커넥터의 라이선스 (License of the RabbitMQ Connector)
Flink의 RabbitMQ 커넥터는 "RabbitMQ AMQP Java Client"에 대한 Maven 의존성을 정의하며, Mozilla Public License 1.1("MPL"), GNU General Public License version 2("GPL"), Apache License version 2("ASL")의 세 가지 라이선스가 적용돼요.
Flink 자체는 "RabbitMQ AMQP Java Client"의 소스 코드를 재사용하지도 않고, 그 바이너리를 패키징하지도 않아요.
Flink의 RabbitMQ 커넥터를 기반으로 파생 작업을 만들어 배포하는 사용자(그에 따라 "RabbitMQ AMQP Java Client"를 재배포)는 Mozilla Public License 1.1("MPL"), GNU General Public License version 2("GPL"), Apache License version 2("ASL")에 선언된 조건의 적용을 받을 수 있다는 점을 인지해야 해요.
RabbitMQ 커넥터
이 커넥터는 RabbitMQ에서 데이터 스트림에 접근할 수 있게 해줘요. 이 커넥터를 사용하려면 프로젝트에 다음 의존성을 추가해요.
현재 Flink 2.3 버전용 커넥터는 아직 제공되지 않아요.
PyFlink 작업에서 사용하려면 다음 의존성이 필요해요:
| 버전(Version) | PyFlink JAR |
|---|---|
| flink-connector-rabbitmq | Flink 2.3 버전용 SQL jar는 아직 없어요. |
PyFlink에서 JAR을 사용하는 방법에 대한 자세한 내용은 Python dependency management를 참조하세요.
스트리밍 커넥터는 현재 바이너리 배포판의 일부가 아니라는 점에 유의하세요. 클러스터 실행을 위한 링크 방법은 여기를 참조하세요.
RabbitMQ 설치 (Installing RabbitMQ)
RabbitMQ 다운로드 페이지의 지침을 따르세요. 설치 후 서버가 자동으로 시작되며, RabbitMQ에 연결하는 애플리케이션을 실행할 수 있어요.
RabbitMQ 소스 (RabbitMQ Source)
이 커넥터는 RabbitMQ 큐에서 메시지를 소비하기 위한 RMQSource 클래스를 제공해요. 이 소스는 Flink로 어떻게 구성되는지에 따라 세 가지 수준의 보장을 제공해요:
-
정확히 한 번 (Exactly-once): RabbitMQ 소스로 정확히 한 번 보장을 달성하려면 다음이 필요해요 -
- 체크포인팅 활성화: 체크포인팅이 활성화되면, 체크포인트가 완료될 때만 메시지가 승인되어(RabbitMQ 큐에서 제거)요.
- 상관 ID(correlation ids) 사용: 상관 ID는 RabbitMQ 애플리케이션 기능이에요. RabbitMQ에 메시지를 주입할 때 메시지 속성에 설정해야 해요. 상관 ID는 체크포인트에서 복원할 때 재처리된 메시지를 소스가 중복 제거하는 데 사용돼요.
- 비-병렬 소스 (Non-parallel source): 정확히 한 번을 달성하려면 소스가 비-병렬(병렬도 1)이어야 해요. 이 제한은 주로 단일 큐에서 여러 소비자로 메시지를 분배하는 RabbitMQ의 방식 때문이에요.
-
최소 한 번 (At-least-once): 체크포인팅이 활성화되어 있지만 상관 ID를 사용하지 않거나 소스가 병렬이면, 소스는 최소 한 번 보장만 제공해요.
-
보장 없음 (No guarantee): 체크포인팅이 활성화되지 않으면, 소스는 강력한 전달 보장이 없어요. 이 설정에서 소스는 Flink의 체크포인팅과 협력하는 대신, 소스가 메시지를 받아 처리하는 즉시 자동으로 승인돼요.
다음은 정확히 한 번 RabbitMQ 소스를 설정하는 코드 예시예요. 인라인 주석은 더 완화된 보장을 위해 구성의 어떤 부분을 무시할 수 있는지 설명해요.
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// checkpointing is required for exactly-once or at-least-once guarantees
env.enableCheckpointing(...);
final RMQConnectionConfig connectionConfig = new RMQConnectionConfig.Builder()
.setHost("localhost")
.setPort(5000)
...
.build();
final DataStream<String> stream = env
.addSource(new RMQSource<String>(
connectionConfig, // config for the RabbitMQ connection
"queueName", // name of the RabbitMQ queue to consume
true, // use correlation ids; can be false if only at-least-once is required
new SimpleStringSchema())) // deserialization schema to turn messages into Java objects
.setParallelism(1); // non-parallel source is only required for exactly-once
val env = StreamExecutionEnvironment.getExecutionEnvironment
// checkpointing is required for exactly-once or at-least-once guarantees
env.enableCheckpointing(...)
val connectionConfig = new RMQConnectionConfig.Builder()
.setHost("localhost")
.setPort(5000)
...
.build
val stream = env
.addSource(new RMQSource[String](
connectionConfig, // config for the RabbitMQ connection
"queueName", // name of the RabbitMQ queue to consume
true, // use correlation ids; can be false if only at-least-once is required
new SimpleStringSchema)) // deserialization schema to turn messages into Java objects
.setParallelism(1) // non-parallel source is only required for exactly-once
env = StreamExecutionEnvironment.get_execution_environment()
# checkpointing is required for exactly-once or at-least-once guarantees
env.enable_checkpointing(...)
connection_config = RMQConnectionConfig.Builder() \
.set_host("localhost") \
.set_port(5000) \
...
.build()
stream = env \
.add_source(RMQSource(
connection_config,
"queueName",
True,
SimpleStringSchema(),
)) \
.set_parallelism(1)
서비스 품질 (QoS) / 소비자 프리페치 (Consumer Prefetch)
RabbitMQ 소스는 RMQConnectionConfig를 통해 소스의 채널에 basicQos를 설정하는 간단한 방법을 제공해요.
병렬 소스마다 연결/채널이 하나씩 있으므로, 이 프리페치 수는 소스의 병렬도와 곱해져 한 번에 작업에 전송될 수 있는 총 미승인(unacknowledged) 메시지 수가 돼요.
더 복잡한 구성이 필요하면 RMQSource#setupChannel(Connection)을 오버라이드해 수동으로 구성할 수 있어요.
final RMQConnectionConfig connectionConfig = new RMQConnectionConfig.Builder()
.setPrefetchCount(30_000)
...
.build();
val connectionConfig = new RMQConnectionConfig.Builder()
.setPrefetchCount(30000)
...
.build
connection_config = RMQConnectionConfig.Builder() \
.set_prefetch_count(30000) \
...
.build()
프리페치 수는 기본적으로 설정되지 않으므로, RabbitMQ 서버가 무제한 메시지를 보낸다는 뜻이에요. 프로덕션에서는 이 값을 설정하는 것이 가장 좋아요. 높은 볼륨의 큐와 체크포인팅이 활성화된 경우, 메시지는 활성화되면 체크포인트에서만 승인되므로 낭비되는 사이클을 줄이기 위해 일부 튜닝이 필요할 수 있어요.
QoS와 프리페치에 대한 자세한 내용은 여기에서, AMQP 0-9-1에서 사용 가능한 옵션에 대한 자세한 내용은 여기에서 확인할 수 있어요.
RabbitMQ 싱크 (RabbitMQ Sink)
이 커넥터는 RabbitMQ 큐에 메시지를 보내기 위한 RMQSink 클래스를 제공해요. 다음은 RabbitMQ 싱크를 설정하는 코드 예시예요.
final DataStream<String> stream = ...
final RMQConnectionConfig connectionConfig = new RMQConnectionConfig.Builder()
.setHost("localhost")
.setPort(5000)
...
.build();
stream.addSink(new RMQSink<String>(
connectionConfig, // config for the RabbitMQ connection
"queueName", // name of the RabbitMQ queue to send messages to
new SimpleStringSchema())); // serialization schema to turn Java objects to messages
val stream: DataStream[String] = ...
val connectionConfig = new RMQConnectionConfig.Builder()
.setHost("localhost")
.setPort(5000)
...
.build
stream.addSink(new RMQSink[String](
connectionConfig, // config for the RabbitMQ connection
"queueName", // name of the RabbitMQ queue to send messages to
new SimpleStringSchema)) // serialization schema to turn Java objects to messages
stream = ...
connection_config = RMQConnectionConfig.Builder() \
.set_host("localhost") \
.set_port(5000) \
...
.build()
stream.add_sink(RMQSink(
connection_config, # config for the RabbitMQ connection
'queueName', # name of the RabbitMQ queue to send messages to
SimpleStringSchema())) # serialization schema to turn Java objects to messages
RabbitMQ에 대한 자세한 내용은 여기에서 확인할 수 있어요.