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 积分,解读结果公开显示在下面。
还没有人解读过这个文件。
