Java消息接收案例如何实现

wen java案例 23

本文目录导读:

Java消息接收案例如何实现

  1. TCP Socket 消息接收(基础案例)
  2. RabbitMQ 消息接收(企业级案例)
  3. Redis 消息订阅(发布/订阅模式)
  4. Kafka 消息消费(大数据处理)
  5. WebSocket 消息接收(实时通信)
  6. 综合示例:消息接收框架
  7. 选择建议

我来为你详细介绍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();
    }
}

选择建议

  1. 简单测试:使用TCP Socket
  2. 企业级消息:使用RabbitMQ
  3. 实时通信:使用WebSocket
  4. 大数据处理:使用Kafka
  5. 简单发布订阅:使用Redis

根据你的具体需求选择合适的消息接收方案。

抱歉,评论功能暂时关闭!