Apache Pulsar消息传递语义
创始人
2024-09-04 19:01:31
0

在Apache Pulsar中,消息传递语义包括“至少一次”(at least once)和“最多一次”(at most once)语义。

  1. 至少一次(At least once)语义: 在至少一次语义中,消息保证至少会被传递一次。这可以通过设置消息的确认模式为ACKNOWLEDGE或AUTO_ACKNOWLEDGE来实现。

示例代码:

import org.apache.pulsar.client.api.*;
import org.apache.pulsar.client.api.schema.GenericRecord;
import org.apache.pulsar.client.api.schema.GenericSchema;

public class AtLeastOnceExample {
    public static void main(String[] args) throws PulsarClientException {
        String serviceUrl = "pulsar://localhost:6650";
        String topic = "persistent://public/default/my-topic";

        PulsarClient client = PulsarClient.builder()
                .serviceUrl(serviceUrl)
                .build();

        Consumer consumer = client.newConsumer(GenericSchema.of())
                .topic(topic)
                .subscriptionName("my-subscription")
                .subscriptionType(SubscriptionType.Shared)
                .ackTimeout(0, TimeUnit.SECONDS)
                .subscribe();

        while (true) {
            Message message = consumer.receive();
            try {
                System.out.println("Received message: " + message.getValue());
                // 处理消息的逻辑

                consumer.acknowledge(message);
            } catch (Exception e) {
                // 处理异常情况
                consumer.negativeAcknowledge(message);
            }
        }
    }
}
  1. 最多一次(At most once)语义: 在最多一次语义中,消息可能会丢失或重复传递。这可以通过设置消息的确认模式为AUTO_ACKNOWLEDGE或AUTO_ACKNOWLEDGE_ASYNC来实现。

示例代码:

import org.apache.pulsar.client.api.*;
import org.apache.pulsar.client.api.schema.GenericRecord;
import org.apache.pulsar.client.api.schema.GenericSchema;

public class AtMostOnceExample {
    public static void main(String[] args) throws PulsarClientException {
        String serviceUrl = "pulsar://localhost:6650";
        String topic = "persistent://public/default/my-topic";

        PulsarClient client = PulsarClient.builder()
                .serviceUrl(serviceUrl)
                .build();

        Consumer consumer = client.newConsumer(GenericSchema.of())
                .topic(topic)
                .subscriptionName("my-subscription")
                .subscriptionType(SubscriptionType.Shared)
                .ackTimeout(0, TimeUnit.SECONDS)
                .subscribe();

        while (true) {
            Message message = consumer.receive();
            System.out.println("Received message: " + message.getValue());
            // 处理消息的逻辑
        }
    }
}

请注意,以上示例代码仅展示了如何设置消息传递语义,实际使用中还需要根据具体的业务逻辑进行相应的处理和错误处理。

相关内容

热门资讯

安卓换鸿蒙系统会卡吗,体验流畅... 最近手机圈可是热闹非凡呢!不少安卓用户都在议论纷纷,说鸿蒙系统要来啦!那么,安卓手机换上鸿蒙系统后,...
app安卓系统登录不了,解锁登... 最近是不是你也遇到了这样的烦恼:手机里那个心爱的APP,突然就登录不上了?别急,让我来帮你一步步排查...
安卓系统拦截短信在哪,安卓系统... 你是不是也遇到了这种情况:手机里突然冒出了很多垃圾短信,烦不胜烦?别急,今天就来教你怎么在安卓系统里...
安卓系统要维护多久,安卓系统维... 你有没有想过,你的安卓手机里那个陪伴你度过了无数日夜的安卓系统,它究竟要陪伴你多久呢?这个问题,估计...
windows官网系统多少钱 Windows官网系统价格一览:了解正版Windows的购买成本Windows 11官方价格解析微软...
安卓系统如何卸载app,轻松掌... 手机里的App越来越多,是不是感觉内存不够用了?别急,今天就来教你怎么轻松卸载安卓系统里的App,让...
怎么复制照片安卓系统,操作步骤... 亲爱的手机控们,是不是有时候想把自己的手机照片分享给朋友,或者备份到电脑上呢?别急,今天就来教你怎么...
安卓系统应用怎么重装,安卓应用... 手机里的安卓应用突然罢工了,是不是让你头疼不已?别急,今天就来手把手教你如何重装安卓系统应用,让你的...
iwatch怎么连接安卓系统,... 你有没有想过,那款时尚又实用的iWatch,竟然只能和iPhone好上好?别急,今天就来给你揭秘,怎...
iphone系统与安卓系统更新... 最近是不是你也遇到了这样的烦恼?手机更新系统总是失败,急得你团团转。别急,今天就来给你揭秘为什么iP...