如何在Apache Flink中连接RabbitMQ Stream而非队列
在Apache Flink中连接RabbitMQ Stream的实现方案
问题:是否可以在Apache Flink中连接RabbitMQ Stream而非队列?我已成功将RabbitMQ队列作为数据源在Apache Flink中运行,但不知如何实现RabbitMQ Stream的连接。
解决方案(更新)
我已通过实现RichParallelSourceFunction并完成run()方法的逻辑编写解决该问题,核心步骤如下:
- 实例化连接RabbitMQ Stream的环境
- 通过
SourceContext收集从RabbitMQ Stream接收的所有消息
完整实现代码
public class RabbitMqSource2 extends RichParallelSourceFunction<String> { private volatile boolean running = true; private Environment environment; private final String streamName = "superStreamName"; private final String consumerName = "Consumer1"; @Override public void run(SourceContext<String> sourceContext) throws Exception { if (Objects.isNull(environment)) { environment = Environment.builder().uris(Arrays.asList( "rabbitmq-stream://usrname:pwd1@Server1:5522/%2f", "rabbitmq-stream://usrname:pwd1@Server2:5522/%2f", "rabbitmq-stream://usrname:pwd1@Server3:5522/%2f") ) .host("Server3") .port(5522) .username("usrname") .password("pwd1") .virtualHost("/") .build(); } environment.consumerBuilder() .superStream(streamName) .name(consumerName) .singleActiveConsumer() .messageHandler((context, message) -> { if(message != null && message.getProperties() != null && message.getProperties().getMessageIdAsString() != null) { String msg = new String(message.getBodyAsBinary()); sourceContext.collect(msg); } }) .build(); } @Override public void cancel() { environment.close(); } }
内容的提问来源于stack exchange,提问作者Nanda
相关产品推荐
相关产品推荐

