使用 reader 接口, Pulsar客戶(hù)可以在主題中“手動(dòng)定位”自己,從指定的消息開(kāi)始向前讀取所有消息
下面是一個(gè)示例:
創(chuàng)新互聯(lián)堅(jiān)持“要么做到,要么別承諾”的工作理念,服務(wù)領(lǐng)域包括:網(wǎng)站設(shè)計(jì)制作、成都網(wǎng)站設(shè)計(jì)、企業(yè)官網(wǎng)、英文網(wǎng)站、手機(jī)端網(wǎng)站、網(wǎng)站推廣等服務(wù),滿(mǎn)足客戶(hù)于互聯(lián)網(wǎng)時(shí)代的吉木薩爾網(wǎng)站設(shè)計(jì)、移動(dòng)媒體設(shè)計(jì)的需求,幫助企業(yè)找到有效的互聯(lián)網(wǎng)解決方案。努力成為您成熟可靠的網(wǎng)絡(luò)建設(shè)合作伙伴!
import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.Reader;
import org.apache.pulsar.client.impl.schema.JSONSchema;
public class ReaderTest{
public static void main(String[] args) {
String url = "http://192.168.1.48:8080";
try{
PulsarClient client =PulsarClient.builder()
.serviceUrl(url)
.build();
Reader reader=client.newReader(JSONSchema.of(UserModel.class))
.topic("my-tenant/my-namespace/testschema-topic")
.startMessageId(MessageId.earliest) //MessageId.earliest最早 MessageId.latest 最新 MessageId斷點(diǎn)
.create();
while (true) {
Message userModelmsg = reader.readNext();
UserModel userModel=userModelmsg.getValue();//業(yè)務(wù)數(shù)據(jù)
MessageId messageId=userModelmsg.getMessageId();//斷點(diǎn)
System.out.println("receive message: " +userModel.getName()+"="+userModel.getAge()+"="+messageId.toString());
}
}catch(Exception e){
e.printStackTrace();
}
}
}