本文目录导读:

- TCP Socket 消息接收(基础案例)
- RabbitMQ 消息接收(企业级案例)
- Redis 消息订阅(发布/订阅模式)
- Kafka 消息消费(大数据处理)
- WebSocket 消息接收(实时通信)
- 综合示例:消息接收框架
- 选择建议
我来为你详细介绍Java消息接收的几种实现方式,从简单到复杂:
TCP Socket 消息接收(基础案例)
服务端代码
import java.io.*;
import java.net.*;
public class TcpServer {
private ServerSocket serverSocket;
public TcpServer(int port) throws IOException {
serverSocket = new ServerSocket(port);
System.out.println("TCP服务器启动,监听端口: " + port);
}
public void start() {
while (true) {
try {
// 接收客户端连接
Socket clientSocket = serverSocket.accept();
System.out.println("客户端连接: " + clientSocket.getInetAddress());
// 为每个客户端创建新线程处理
new Thread(new ClientHandler(clientSocket)).start();
} catch (IOException e) {
e.printStackTrace();
}
}
}
// 客户端处理线程
static class ClientHandler implements Runnable {
private Socket socket;
public ClientHandler(Socket socket) {
this.socket = socket;
}
@Override
public void run() {
try (BufferedReader reader = new BufferedReader(
new InputStreamReader(socket.getInputStream()));
PrintWriter writer = new PrintWriter(
socket.getOutputStream(), true)) {
String message;
while ((message = reader.readLine()) != null) {
System.out.println("收到消息: " + message);
// 发送响应
writer.println("服务端已收到: " + message);
}
} catch (IOException e) {
e.printStackTrace();
} finally {
try {
socket.close();
} catch (IOException e) {
e.printStackTrace();
}
}
}
}
public static void main(String[] args) throws IOException {
TcpServer server = new TcpServer(8080);
server.start();
}
}
客户端代码
import java.io.*;
import java.net.*;
public class TcpClient {
public static void main(String[] args) {
String host = "localhost";
int port = 8080;
try (Socket socket = new Socket(host, port);
PrintWriter writer = new PrintWriter(
socket.getOutputStream(), true);
BufferedReader reader = new BufferedReader(
new InputStreamReader(socket.getInputStream()));
BufferedReader consoleReader = new BufferedReader(
new InputStreamReader(System.in))) {
System.out.println("已连接到服务器");
String userInput;
while ((userInput = consoleReader.readLine()) != null) {
// 发送消息
writer.println(userInput);
// 接收响应
String response = reader.readLine();
System.out.println("服务器响应: " + response);
}
} catch (IOException e) {
e.printStackTrace();
}
}
}
RabbitMQ 消息接收(企业级案例)
Maven依赖
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.16.0</version>
</dependency>
消息接收者
import com.rabbitmq.client.*;
public class RabbitMQConsumer {
private final static String QUEUE_NAME = "hello";
public static void main(String[] args) throws Exception {
// 创建连接工厂
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setUsername("guest");
factory.setPassword("guest");
// 创建连接和通道
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// 声明队列
channel.queueDeclare(QUEUE_NAME, false, false, false, null);
System.out.println("等待消息...");
// 设置消息接收回调
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println("收到消息: " + message);
// 手动确认消息
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
};
// 取消回调
CancelCallback cancelCallback = consumerTag -> {
System.out.println("消费被取消");
};
// 消费消息
channel.basicConsume(QUEUE_NAME, false, deliverCallback, cancelCallback);
// 保持程序运行
Thread.sleep(Long.MAX_VALUE);
}
}
}
Redis 消息订阅(发布/订阅模式)
Maven依赖
<dependency>
<groupId>redis.clients</groupId>
<artifactId>jedis</artifactId>
<version>4.3.1</version>
</dependency>
消息监听器
import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisPubSub;
public class RedisSubscriber {
public static void main(String[] args) {
// 创建Jedis连接
Jedis jedis = new Jedis("localhost", 6379);
// 创建订阅者
JedisPubSub subscriber = new JedisPubSub() {
@Override
public void onMessage(String channel, String message) {
System.out.println("频道: " + channel + ", 消息: " + message);
}
@Override
public void onSubscribe(String channel, int subscribedChannels) {
System.out.println("订阅频道: " + channel);
}
@Override
public void onUnsubscribe(String channel, int subscribedChannels) {
System.out.println("取消订阅: " + channel);
}
};
// 订阅频道
new Thread(() -> {
jedis.subscribe(subscriber, "news", "sports");
}).start();
// 保持程序运行
try {
Thread.sleep(Long.MAX_VALUE);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
Kafka 消息消费(大数据处理)
Maven依赖
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.4.0</version>
</dependency>
Kafka消费者
import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
public class KafkaConsumer {
public static void main(String[] args) {
// 配置消费者属性
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
// 创建消费者
try (Consumer<String, String> consumer =
new KafkaConsumer<>(props)) {
// 订阅主题
consumer.subscribe(Collections.singletonList("my-topic"));
System.out.println("开始消费消息...");
while (true) {
// 拉取消息
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("分区: %d, 偏移量: %d, 键: %s, 值: %s%n",
record.partition(),
record.offset(),
record.key(),
record.value());
}
// 提交偏移量
consumer.commitAsync();
}
}
}
}
WebSocket 消息接收(实时通信)
Spring Boot WebSocket 配置
import org.springframework.context.annotation.Configuration;
import org.springframework.web.socket.config.annotation.EnableWebSocket;
import org.springframework.web.socket.config.annotation.WebSocketConfigurer;
import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry;
@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {
@Override
public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
registry.addHandler(new MyWebSocketHandler(), "/ws")
.setAllowedOrigins("*");
}
}
WebSocket处理器
import org.springframework.web.socket.*;
import org.springframework.web.socket.handler.TextWebSocketHandler;
public class MyWebSocketHandler extends TextWebSocketHandler {
@Override
public void afterConnectionEstablished(WebSocketSession session) {
System.out.println("新连接: " + session.getId());
}
@Override
protected void handleTextMessage(WebSocketSession session,
TextMessage message) throws Exception {
String payload = message.getPayload();
System.out.println("收到消息: " + payload);
// 回复消息
session.sendMessage(new TextMessage("已收到: " + payload));
}
@Override
public void afterConnectionClosed(WebSocketSession session,
CloseStatus status) {
System.out.println("连接关闭: " + session.getId());
}
}
综合示例:消息接收框架
import java.util.function.Consumer;
// 消息接收器接口
interface MessageReceiver {
void start();
void stop();
void subscribe(String topic, Consumer<String> handler);
void unsubscribe(String topic);
}
// TCP实现
class TcpMessageReceiver implements MessageReceiver {
private ServerSocket serverSocket;
private Map<String, Consumer<String>> handlers = new HashMap<>();
private volatile boolean running = false;
@Override
public void start() {
running = true;
// 启动接收线程
}
@Override
public void stop() {
running = false;
// 关闭连接
}
@Override
public void subscribe(String topic, Consumer<String> handler) {
handlers.put(topic, handler);
}
@Override
public void unsubscribe(String topic) {
handlers.remove(topic);
}
}
// 使用示例
public class MessageReceiverDemo {
public static void main(String[] args) {
MessageReceiver receiver = new TcpMessageReceiver();
// 订阅消息
receiver.subscribe("orders", message -> {
System.out.println("订单消息: " + message);
// 处理订单
});
receiver.subscribe("notifications", message -> {
System.out.println("通知消息: " + message);
// 发送通知
});
// 启动接收器
receiver.start();
}
}
选择建议
- 简单测试:使用TCP Socket
- 企业级消息:使用RabbitMQ
- 实时通信:使用WebSocket
- 大数据处理:使用Kafka
- 简单发布订阅:使用Redis
根据你的具体需求选择合适的消息接收方案。