您好,登錄后才能下訂單哦!
import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.SubscriptionInitialPosition;
import org.apache.pulsar.client.api.SubscriptionType;
import org.apache.pulsar.client.impl.schema.JSONSchema;
public class ReceiveMsgTest {
public static void main(String[] args) {
String url = "http://192.168.1.48:8080";
try{
PulsarClient client =PulsarClient.builder()
.serviceUrl(url)
.build();
Consumer<UserModel> consumer=client.newConsumer(JSONSchema.of(UserModel.class))
.topic("my-tenant/my-namespace/testschema-topic")
.subscriptionType(SubscriptionType.Exclusive)//訂閱模式 Exclusive(獨占,默認模式) Failover(災備)Shared(共享)
.subscriptionName("wbq_1")//訂閱者名稱
.subscribe();
while (true) {
Message<UserModel> userModelmsg = consumer.receive();
UserModel userModel=userModelmsg.getValue();
System.out.println("receive message: " +userModel.getName()+"="+userModel.getAge());
consumer.acknowledge(userModelmsg.getMessageId());//應答后此訂閱者不會在收到此消息
}
}catch(Exception e){
e.printStackTrace();
}
}
}
public class UserModel {
private String name;
private int age;
public String getName() {
return name;
}
public void setName(String name) {
this.name = name;
}
public int getAge() {
return age;
}
public void setAge(int age) {
this.age = age;
}
}
免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。