oasys-service/oasys-chat/src/main/java/cn/linter/oasys/chat/config/PulsarConfig.java

(开头部分) 1KB

这里只显示每个文件的开头 60 行。登录后可以解锁完整代码。

package cn.linter.oasys.chat.config;

import cn.linter.oasys.chat.handler.ChatMessageListener;
import lombok.extern.slf4j.Slf4j;
import org.apache.pulsar.client.api.*;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

/**
 * @author wangxiaoyang
 */
@Slf4j
@Configuration
public class PulsarConfig {

    @Value("${pulsar.host}")
    private String pulsarHost;
    @Value("${pulsar.topic}")
    private String pulsarTopic;
    @Value("${pulsar.subscription}")
    private String pulsarSubscription;

    @Bean
    public PulsarClient pulsarClient() {
        try {
            return PulsarClient.builder()
                    .serviceUrl("pulsar://" + pulsarHost)
                    .build();
        } catch (PulsarClientException e) {
            log.error("Pulsar连接失败", e);
            return null;
        }
    }

    @Bean
    public Producer<String> producer(PulsarClient pulsarClient) {
        try {
            return pulsarClient.newProducer(Schema.STRING)
                    .topic(pulsarTopic)
                    .create();
        } catch (PulsarClientException e) {
            log.error("Pulsar生产者创建失败", e);
            return null;
        }
    }

    @Bean
    public Consumer<String> consumer(PulsarClient pulsarClient, ChatMessageListener messageListener) {
        try {
            return pulsarClient.newConsumer(Schema.STRING)
                    .topic(pulsarTopic)
                    .subscriptionName(pulsarSubscription)
                    .messageListener(messageListener)
                    .subscribe();
        } catch (PulsarClientException e) {
            log.error("Pulsar生产者创建失败", e);
            return null;
        }
    }
后面还有 2 行代码,购买后查看完整代码

24 小时内免费解锁 3 个项目,之后 1 积分/个。 规则说明

AI 解读

登录后可用,每次 10 积分,解读结果公开显示在下面。

还没有人解读过这个文件。