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

Hazelcast服务端向客户端发送Topic消息失败,求示例程序

Hazelcast ITopic 服务端向客户端发送消息示例

我来帮你搞定这个问题!你遇到的情况大概率是客户端的订阅逻辑没处理到位,或者集群配置不一致导致的。下面给你一套完整的服务端和客户端示例代码,亲测可以实现服务端主动推消息给客户端:

服务端代码

import com.hazelcast.core.Hazelcast;
import com.hazelcast.core.HazelcastInstance;
import com.hazelcast.core.ITopic;

public class TopicServer {
    public static void main(String[] args) {
        // 初始化Hazelcast实例,默认会自动组成集群
        HazelcastInstance hazelcastInstance = Hazelcast.newHazelcastInstance();
        
        // 获取指定名称的Topic,注意名称要和客户端完全一致
        ITopic<String> topic = hazelcastInstance.getTopic("server-to-client-topic");
        
        // 模拟服务端主动发送消息,这里可以替换成你的业务触发逻辑
        for (int i = 0; i < 5; i++) {
            String message = "服务端推送消息:" + i;
            topic.publish(message);
            System.out.println("服务端已发送:" + message);
            
            try {
                Thread.sleep(2000); // 每2秒发一条,方便测试
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }
        
        // 注意:实际业务中不要随便关闭实例,这里只是示例结束后清理
        // hazelcastInstance.shutdown();
    }
}

客户端代码

import com.hazelcast.core.Hazelcast;
import com.hazelcast.core.HazelcastInstance;
import com.hazelcast.core.ITopic;
import com.hazelcast.core.MessageListener;

public class TopicClient {
    public static void main(String[] args) {
        // 初始化Hazelcast客户端实例,默认会自动发现集群中的服务端
        HazelcastInstance hazelcastInstance = Hazelcast.newHazelcastInstance();
        
        // 获取和服务端同名的Topic
        ITopic<String> topic = hazelcastInstance.getTopic("server-to-client-topic");
        
        // 注册消息监听器,这是客户端能收到消息的关键!
        topic.addMessageListener(new MessageListener<String>() {
            @Override
            public void onMessage(com.hazelcast.core.Message<String> message) {
                System.out.println("客户端收到消息:" + message.getMessageObject());
            }
        });
        
        // 保持客户端运行,否则程序退出就收不到消息了
        System.out.println("客户端已启动,等待接收消息...");
        try {
            Thread.currentThread().join();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
        
        // 实际业务中按需关闭
        // hazelcastInstance.shutdown();
    }
}

关键注意事项

  • 依赖配置:确保服务端和客户端都引入了相同版本的Hazelcast依赖,比如Maven依赖:
    <dependency>
        <groupId>com.hazelcast</groupId>
        <artifactId>hazelcast</artifactId>
        <version>5.3.6</version> <!-- 用最新稳定版即可 -->
    </dependency>
    
  • 集群一致性:服务端和客户端的Hazelcast配置要一致(比如网络配置、群组名称等),默认配置下只要在同一局域网就能自动发现,不需要额外配置。
  • Topic名称一致:服务端和客户端使用的Topic名称必须完全相同,大小写敏感!
  • 客户端保持运行:客户端程序不能启动后立即退出,否则监听器无法生效,示例中用Thread.currentThread().join()让程序一直运行。
  • 监听器注册时机:客户端要在服务端发送消息前完成监听器注册,否则会错过之前的消息(ITopic默认不保留历史消息,需要持久化的话可以用ReliableTopic)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:33:12