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

Flink 1.3.2中无法启动RabbitMQ Source问题求助

从你贴的错误栈来看,核心问题就是java.net.ConnectException: Connection refused——说白了就是Flink程序压根连不上你的RabbitMQ服务器,咱们一步步来排查解决:

1. 先把端口改对!

你代码里设置的端口是50000,但RabbitMQ的默认AMQP通信端口是5672哦!除非你之前特意修改过RabbitMQ的配置文件改了端口,否则这个端口肯定是错的,先把它改成5672试试:

RMQConnectionConfig connectionConfig = new RMQConnectionConfig.Builder()
    .setHost("localhost")
    .setPort(5672) // 这里改成RabbitMQ默认端口
    .setUserName("root")
    .setPassword("root")
    .setVirtualHost("/").build();

2. 聊聊用户名密码的坑

RabbitMQ默认自带的管理员账号是guest/guest,不过这个账号有个限制:默认只能从本地localhost访问。如果你是在本地同时跑Flink和RabbitMQ,用这个账号完全没问题;但如果你之前自己创建了root用户,那得确认几点:

  • 这个root用户真的在RabbitMQ里存在(可以登录RabbitMQ控制台,去Users页面查)
  • 它拥有访问/这个虚拟主机的权限(去Permissions页面看权限配置)
  • 它对你要消费的queue队列有读写权限

3. 额外的验证小技巧

如果改完端口还是不行,试试这些方法排查:

  • 用telnet localhost 5672或者nc -zv localhost 5672测试端口能不能连通,能通的话说明网络层面没问题
  • 登录RabbitMQ控制台的Connections页面,看看有没有失败的连接记录,里面会有更详细的错误信息,帮你定位是权限还是配置问题
  • 检查Flink依赖的RabbitMQ客户端版本和服务器版本是否兼容,版本差太多也可能导致连接失败

最后给你个改好的示例代码(用默认端口和guest账号):

public class rabbitmq_source {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment envrionment = StreamExecutionEnvironment.getExecutionEnvironment();
        RMQConnectionConfig connectionConfig = new RMQConnectionConfig.Builder()
            .setHost("localhost")
            .setPort(5672) // 默认AMQP端口
            .setUserName("guest")
            .setPassword("guest")
            .setVirtualHost("/").build();
        DataStream<String> stream = envrionment
            .addSource(new RMQSource<String>(
                connectionConfig, 
                "queue", 
                new SimpleStringSchema()));
        stream.print();
        envrionment.execute();
    }
}

内容的提问来源于stack exchange,提问作者Amarjit Dhillon

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 03:44:22