Flink 1.3.2中无法启动RabbitMQ Source问题求助
搞定Flink 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
相关产品推荐
相关产品推荐

