You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 16:32:37