如何在Apache Ignite队列上实现生产者-消费者模型?
Ignite队列事件监听与生产者-消费者模型实现
Ignite队列是否支持原生事件回调?
Apache Ignite 提供的分布式队列 IgniteQueue 本身没有内置类似缓存连续查询的onAdd/onRemove原生回调接口,也没有独立的事件触发机制。但可以结合Ignite现有特性,低成本实现完全等效的监听效果,同时也支持标准的生产者-消费者模式落地。
标准生产者-消费者实现(官方推荐、生产环境首选)
IgniteQueue 本身扩展了JDK BlockingQueue 接口,天然支持阻塞式消费,这种实现性能最高、逻辑最简单,是绝大多数场景的首选:
生产者示例
// 获取或创建分布式队列,参数依次为队列名、容量上限(0代表无界)、队列配置 // setBackups(1) 表示设置1个备份,避免节点宕机丢数据 IgniteQueue<YourBusinessObj> queue = ignite.queue( "biz-queue", 0, new CollectionConfiguration().setBackups(1) ); // 写入数据,队列满时会自动阻塞 queue.put(new YourBusinessObj()); // 也可使用offer指定超时时间 // queue.offer(new YourBusinessObj(), 5, TimeUnit.SECONDS);
消费者示例(等效触发onAdd逻辑)
// 启动消费线程 new Thread(() -> { while (!Thread.currentThread().isInterrupted()) { try { // 队列无数据时自动阻塞,有新元素时立刻返回,等效于onAdd事件触发 YourBusinessObj obj = queue.take(); // 此处执行业务消费逻辑 processBizObj(obj); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } }, "queue-consumer-thread").start();
该方案天然支持多消费者并发消费,Ignite会自动做消费负载均衡,无额外的事件转发开销,性能远高于自定义事件监听方案。
自定义事件监听实现(满足queue.onAdd式回调需求)
如果你确实需要事件回调的编码形式,可以利用Ignite队列底层基于分布式缓存存储的特性,通过缓存事件监听实现等效能力:
Ignite队列对应的底层缓存命名规则为 ignite-sys-cache-queue.{你的队列名},只需要监听该缓存的写入/删除事件即可:
String queueName = "biz-queue"; String queueCacheName = "ignite-sys-cache-queue." + queueName; // 注册连续查询监听缓存事件 ContinuousQuery<Long, YourBusinessObj> continuousQry = new ContinuousQuery<>(); continuousQry.setLocalListener(events -> { for (CacheEntryEvent<? extends Long, ? extends YourBusinessObj> event : events) { if (event.getEventType() == EventType.CREATED) { // 新增元素回调,即onAdd逻辑 System.out.println("监听到队列新增元素: " + event.getValue()); processOnAdd(event.getValue()); } else if (event.getEventType() == EventType.REMOVED) { // 元素移除回调,即onRemove逻辑 System.out.println("监听到队列移除元素: " + event.getOldValue()); processOnRemove(event.getOldValue()); } } }); // 启动监听 ignite.cache(queueCacheName).query(continuousQry);
注意:该方案适合需要异步感知队列元素变更的特殊场景,常规生产者消费者需求优先选择阻塞消费方案即可。
内容的提问来源于stack exchange,提问作者Jiayu
相关产品推荐
相关产品推荐

