如何编写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.
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"); } }
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
- Start your local RabbitMQ server (default port 5672, credentials
guest/guest). - Run the Spring Boot app – it will send 3 test messages to the
topic_logsexchange on startup. - Open a terminal and run the Python consumer with binding keys, e.g.:
This will catch all messages with routing keys starting withpython receive.py "user.*"user.. - Spin up another consumer instance with a different binding:
This will only pick up messages with the exactpython receive.py "order.create"order.createrouting 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.*hitsuser.createbut notuser.profile.update)#matches zero or more words (e.g.,user.#catches all user-related messages)
内容的提问来源于stack exchange,提问作者megna

