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

如何编写Spring Boot RabbitMQ消息发送端对接Python接收端?

Got it, let's get your Spring Boot + Python RabbitMQ setup up and running properly. Below is a complete implementation for both the Spring Boot message producer and the finished Python message consumer, aligned with the topic exchange pattern you started with.

Spring Boot RabbitMQ Producer Implementation

1. Add Maven Dependencies

First, include the necessary RabbitMQ starter in your pom.xml:

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>
</dependencies>

2. Configure RabbitMQ Connection

Update your application.yml (or application.properties) with RabbitMQ server details:

spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest

3. Producer Code

Create a configuration class to declare the topic exchange (code-based declaration ensures consistency even if you reset RabbitMQ):

import org.springframework.amqp.core.TopicExchange;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class RabbitMQConfig {
    public static final String TOPIC_EXCHANGE_NAME = "topic_logs";

    @Bean
    public TopicExchange topicExchange() {
        return new TopicExchange(TOPIC_EXCHANGE_NAME);
    }
}

Then create a service to handle sending messages to the exchange with specific routing keys:

import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

@Service
public class MessageProducer {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void sendMessage(String routingKey, String message) {
        rabbitTemplate.convertAndSend(RabbitMQConfig.TOPIC_EXCHANGE_NAME, routingKey, message);
        System.out.println("Sent message: '" + message + "' with routing key: '" + routingKey + "'");
    }
}

You can test this with a command line runner that sends sample messages on startup:

import org.springframework.boot.CommandLineRunner;
import org.springframework.stereotype.Component;

@Component
public class ProducerRunner implements CommandLineRunner {

    @Autowired
    private MessageProducer producer;

    @Override
    public void run(String... args) throws Exception {
        // Send test messages with different routing keys
        producer.sendMessage("user.create", "New user registered: Alice");
        producer.sendMessage("user.update", "User Bob updated their profile");
        producer.sendMessage("order.create", "New order #123 placed");
    }
}
Finished Python Consumer Code

Here's the complete receive.py with the missing binding keys logic and message handling:

#!/usr/bin/env python
import pika
import sys

def main():
    # Connect to local RabbitMQ server
    connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost'))
    channel = connection.channel()

    # Declare the same topic exchange as the Spring Boot producer
    channel.exchange_declare(exchange='topic_logs', exchange_type='topic')

    # Create an exclusive temporary queue (auto-deleted when consumer disconnects)
    result = channel.queue_declare(queue='', exclusive=True)
    queue_name = result.method.queue

    # Get binding keys from command line arguments (e.g., python receive.py "user.*" "order.create")
    binding_keys = sys.argv[1:]
    if not binding_keys:
        sys.stderr.write("Usage: %s [binding_key]...\n" % sys.argv[0])
        sys.exit(1)

    # Bind the queue to the exchange with each specified binding key
    for binding_key in binding_keys:
        channel.queue_bind(
            exchange='topic_logs', queue=queue_name, routing_key=binding_key)
        print("Binding queue %s to exchange %s with routing key %s" % (queue_name, 'topic_logs', binding_key))

    print(' [*] Waiting for messages. To exit press CTRL+C')

    # Define what to do when a message is received
    def callback(ch, method, properties, body):
        print(" [x] Received %r:%r" % (method.routing_key, body.decode()))

    # Start consuming messages from the queue
    channel.basic_consume(
        queue=queue_name, on_message_callback=callback, auto_ack=True)

    channel.start_consuming()

if __name__ == '__main__':
    try:
        main()
    except KeyboardInterrupt:
        print('Interrupted')
        try:
            sys.exit(0)
        except SystemExit:
            exit(0)

How to Test the Setup

  1. Start your local RabbitMQ server (default port 5672, credentials guest/guest).
  2. Run the Spring Boot app – it will send 3 test messages to the topic_logs exchange on startup.
  3. Open a terminal and run the Python consumer with binding keys, e.g.:
    python receive.py "user.*"
    
    This will catch all messages with routing keys starting with user..
  4. Spin up another consumer instance with a different binding:
    python receive.py "order.create"
    
    This will only pick up messages with the exact order.create routing key.

Quick reminders about topic exchanges:

  • Routing keys use dot-separated words (e.g., user.profile.update)
  • Binding keys support wildcards:
    • * matches exactly one word (e.g., user.* hits user.create but not user.profile.update)
    • # matches zero or more words (e.g., user.# catches all user-related messages)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:26:42