pulsar-client-cpp
Consumer.h
1 
19 #ifndef CONSUMER_HPP_
20 #define CONSUMER_HPP_
21 
22 #include <iostream>
23 #include <pulsar/BrokerConsumerStats.h>
24 #include <pulsar/ConsumerConfiguration.h>
25 #pragma GCC visibility push(default)
26 
27 namespace pulsar {
28 class PulsarWrapper;
29 class ConsumerImplBase;
30 class PulsarFriend;
31 typedef std::shared_ptr<ConsumerImplBase> ConsumerImplBasePtr;
35 class Consumer {
36  public:
40  Consumer();
41 
45  const std::string& getTopic() const;
46 
50  const std::string& getSubscriptionName() const;
51 
66 
78  void unsubscribeAsync(ResultCallback callback);
79 
90  Result receive(Message& msg);
91 
100  Result receive(Message& msg, int timeoutMs);
101 
113  void receiveAsync(ReceiveCallback callback);
114 
126  Result acknowledge(const Message& message);
127  Result acknowledge(const MessageId& messageId);
128 
138  void acknowledgeAsync(const Message& message, ResultCallback callback);
139  void acknowledgeAsync(const MessageId& messageID, ResultCallback callback);
140 
158  Result acknowledgeCumulative(const Message& message);
159  Result acknowledgeCumulative(const MessageId& messageId);
160 
171  void acknowledgeCumulativeAsync(const Message& message, ResultCallback callback);
172  void acknowledgeCumulativeAsync(const MessageId& messageId, ResultCallback callback);
173 
174  Result close();
175 
176  void closeAsync(ResultCallback callback);
177 
178  /*
179  * Pause receiving messages via the messageListener, till resumeMessageListener() is called.
180  */
181  Result pauseMessageListener();
182 
183  /*
184  * Resume receiving the messages via the messageListener.
185  * Asynchronously receive all the messages enqueued from time pauseMessageListener() was called.
186  */
187  Result resumeMessageListener();
188 
199 
212  Result getBrokerConsumerStats(BrokerConsumerStats& brokerConsumerStats);
213 
225  void getBrokerConsumerStatsAsync(BrokerConsumerStatsCallback callback);
226 
237  Result seek(const MessageId& msgId);
238 
249  virtual void seekAsync(const MessageId& msgId, ResultCallback callback);
250 
251  private:
252  ConsumerImplBasePtr impl_;
253  explicit Consumer(ConsumerImplBasePtr);
254 
255  friend class PulsarFriend;
256  friend class PulsarWrapper;
257  friend class PartitionedConsumerImpl;
258  friend class MultiTopicsConsumerImpl;
259  friend class ConsumerImpl;
260  friend class ClientImpl;
261  friend class ConsumerTest;
262 };
263 } // namespace pulsar
264 
265 #pragma GCC visibility pop
266 
267 #endif /* CONSUMER_HPP_ */
pulsar::Consumer::receiveAsync
void receiveAsync(ReceiveCallback callback)
pulsar::MessageId
Definition: MessageId.h:33
pulsar::Consumer::seekAsync
virtual void seekAsync(const MessageId &msgId, ResultCallback callback)
pulsar::Consumer::acknowledgeAsync
void acknowledgeAsync(const Message &message, ResultCallback callback)
pulsar::Result
Result
Definition: Result.h:31
pulsar::Consumer::Consumer
Consumer()
pulsar::Consumer::getBrokerConsumerStats
Result getBrokerConsumerStats(BrokerConsumerStats &brokerConsumerStats)
pulsar::Consumer::getTopic
const std::string & getTopic() const
pulsar::Consumer::unsubscribe
Result unsubscribe()
pulsar::Consumer::redeliverUnacknowledgedMessages
void redeliverUnacknowledgedMessages()
pulsar::Consumer::getSubscriptionName
const std::string & getSubscriptionName() const
pulsar::Message
Definition: Message.h:43
pulsar::Consumer::acknowledgeCumulative
Result acknowledgeCumulative(const Message &message)
pulsar::BrokerConsumerStats
Definition: BrokerConsumerStats.h:35
pulsar::Consumer
Definition: Consumer.h:35
pulsar::Consumer::unsubscribeAsync
void unsubscribeAsync(ResultCallback callback)
pulsar::Consumer::acknowledge
Result acknowledge(const Message &message)
pulsar
Definition: Authentication.h:31
pulsar::Consumer::getBrokerConsumerStatsAsync
void getBrokerConsumerStatsAsync(BrokerConsumerStatsCallback callback)
pulsar::Consumer::seek
Result seek(const MessageId &msgId)
pulsar::Consumer::acknowledgeCumulativeAsync
void acknowledgeCumulativeAsync(const Message &message, ResultCallback callback)
pulsar::Consumer::receive
Result receive(Message &msg)
pulsar::ResultCallback
std::function< void(Result result)> ResultCallback
Callback definition for non-data operation.
Definition: ConsumerConfiguration.h:36