Spring Boot 整合 Apache Pulsar 入门指南

未分类2个月前发布 元一软件
93 0
@Autowired

private PulsarTemplate<User> template;



private static final String USER_TOPIC = "user-topic";



public void sendMessageToPulsarTopic(User user) throws PulsarClientException {

    template.send(USER_TOPIC, user);

}

在上面的代码片段中,我们使用 PulsarTemplate 向 Apache Pulsar 的 user-topic topic 发送了一个 User class 对象。

开发工具

5.2、自定义生产者配置

PulsarTemplate 接受 TypedMessageBuilderCustomizer 来配置发送的信息,并接受 ProducerBuilderCustomizer 来定制生产者的属性。

我们可以使用 TypedMessageBuilderCustomizer 来配置消息延迟、在特定时间发送、禁用复制以及提供其他属性:

public void sendMessageToPulsarTopic(User user) throws PulsarClientException {

    template.newMessage(user)

      .withMessageCustomizer(mc -> {

        mc.deliverAfter(10L, TimeUnit.SECONDS);

      })

      .send();

}

ProducerBuilderCustomizer 可用于添加访问模式、自定义消息路由和拦截器,以及启用或禁用分块(chunking)和批处理:

public void sendMessageToPulsarTopic(User user) throws PulsarClientException {

    template.newMessage(user)

      .withProducerCustomizer(pc -> {

        pc.accessMode(ProducerAccessMode.Shared);

      })

      .send();

}

6、消费者

在向 topic 发布消息后,我们现在要为同一个 topic 建立一个 listener。要启用对 topic 的监听,需要用 @PulsarListener 注解 listener 方法。

软件

Spring Boot 会为 listener 方法配置所有必要的组件。

我们还需要使用 @EnablePulsar 注解来启用 PulsarListener。

6.1. 接收消息

首先要为前一节创建的 “string-topic” 创建一个 listener 方法:

@Service

public class PulsarConsumer {



    private static final String STRING_TOPIC = "string-topic";



    @PulsarListener(

      subscriptionName = "string-topic-subscription",

      topics = STRING_TOPIC,

      subscriptionType = SubscriptionType.Shared

    )

    public void stringTopicListener(String str) {

        LOGGER.info("Received String message: {}", str);

    }

}

在 PulsarListener 注解中,我们在 topicName 属性中配置了该方法将监听的 topic,并在 subscriptionName 属性中给出了订阅名称。

编程

现在,让我们为 User 类使用的 user-topic 创建一个 listener 方法:

private static final String USER_TOPIC = "user-topic";



@PulsarListener(

    subscriptionName = "user-topic-subscription",

    topics = USER_TOPIC,

    schemaType = SchemaType.JSON

)

public void userTopicListener(User user) {

    LOGGER.info("Received user object with email: {}", user.getEmail());

}

除了先前的 Listener 方法中提供的属性外,我们还添加了一个 schemaType 属性,其值与生产者中的值相同。

还需要在 main class 上添加 @EnablePulsar 注解:

@EnablePulsar

@SpringBootApplication

public class SpringPulsarApplication {



    public static void main(String[] args) {

        SpringApplication.run(SpringPulsarApplication.class, args);

    }

}

6.2、自定义消费者配置

除订阅名称和 schema type 外,PulsarListener 还可用于配置自动启动、批处理和确认模式等属性:

Java(编程语言)

@PulsarListener(

  subscriptionName = "user-topic-subscription",

  topics = USER_TOPIC,

  subscriptionType = SubscriptionType.Shared,

  schemaType = SchemaType.JSON,

  ackMode = AckMode.RECORD,

  properties = {"ackTimeout=60s"}

)

public void userTopicListener(User user) {

    LOGGER.info("Received user object with email: {}", user.getEmail());

}

在这里,我们将确认模式设置为 Record,并将确认超时设置为 60 秒。

7、使用死信 Topic

如果信息确认超时或服务器接收到 nack,Pulsar 就会尝试重发一定次数的信息。这些重试次数用完后,这些未送达的信息会被发送到称为死信队列(DLQ)的队列中。

此选项仅适用于共享(Shared)订阅类型。要为我们的 user-topic 队列配置 DLQ,我们首先要创建一个 DeadLetterPolicy Bean,它将定义尝试重新交付的次数以及用作 DLQ 的队列名称:

private static final String USER_DEAD_LETTER_TOPIC = "user-dead-letter-topic";
@Bean

DeadLetterPolicy deadLetterPolicy() {

    return DeadLetterPolicy.builder()

      .maxRedeliverCount(10)

      .deadLetterTopic(USER_DEAD_LETTER_TOPIC)

      .build();

}

现在,我们将把该策略添加到之前创建的 PulsarListener 中:

@PulsarListener(

  subscriptionName = "user-topic-subscription",

  topics = USER_TOPIC,

  subscriptionType = SubscriptionType.Shared,

  schemaType = SchemaType.JSON,

  deadLetterPolicy = "deadLetterPolicy",

  properties = {"ackTimeout=60s"}

)

public void userTopicListener(User user) {

    LOGGER.info("Received user object with email: {}", user.getEmail());

}

在这里,我们将 userTopicListener 配置为使用之前创建的 deadLetterPolicy,并将确认时间配置为 60 秒。

技术参考信息

我们可以创建一个单独的 Listener 来处理 DQL 中的信息:

@PulsarListener(

  subscriptionName = "dead-letter-topic-subscription",

  topics = USER_DEAD_LETTER_TOPIC,

  subscriptionType = SubscriptionType.Shared

)

public void userDlqTopicListener(User user) {

    LOGGER.info("Received user object in user-DLQ with email: {}", user.getEmail());

}

8、总结

在本教程中,我们学习了如何在 Spring Boot 应用程序中整合,使用 Apache Pulsar,以及如何更改生产者和消费者的默认配置。


参考:https://www.baeldung.com/spring-boot-apache-pulsar

© 版权声明