librdkafka
The Apache Kafka C/C++ client library
|
Base handle, super class for specific clients. More...
#include <rdkafkacpp.h>
Public Member Functions | |
virtual std::string | name () const =0 |
virtual std::string | memberid () const =0 |
Returns the client's broker-assigned group member id. More... | |
virtual int | poll (int timeout_ms)=0 |
Polls the provided kafka handle for events. More... | |
virtual int | outq_len ()=0 |
Returns the current out queue length. More... | |
virtual ErrorCode | metadata (bool all_topics, const Topic *only_rkt, Metadata **metadatap, int timeout_ms)=0 |
Request Metadata from broker. More... | |
virtual ErrorCode | pause (std::vector< TopicPartition * > &partitions)=0 |
Pause producing or consumption for the provided list of partitions. More... | |
virtual ErrorCode | resume (std::vector< TopicPartition * > &partitions)=0 |
Resume producing or consumption for the provided list of partitions. More... | |
virtual ErrorCode | query_watermark_offsets (const std::string &topic, int32_t partition, int64_t *low, int64_t *high, int timeout_ms)=0 |
Query broker for low (oldest/beginning) and high (newest/end) offsets for partition. More... | |
virtual ErrorCode | get_watermark_offsets (const std::string &topic, int32_t partition, int64_t *low, int64_t *high)=0 |
Get last known low (oldest/beginning) and high (newest/end) offsets for partition. More... | |
virtual ErrorCode | offsetsForTimes (std::vector< TopicPartition * > &offsets, int timeout_ms)=0 |
Look up the offsets for the given partitions by timestamp. More... | |
virtual Queue * | get_partition_queue (const TopicPartition *partition)=0 |
Retrieve queue for a given partition. More... | |
virtual ErrorCode | set_log_queue (Queue *queue)=0 |
Forward librdkafka logs (and debug) to the specified queue for serving with one of the ..poll() calls. More... | |
virtual void | yield ()=0 |
Cancels the current callback dispatcher (Handle::poll(), KafkaConsumer::consume(), etc). More... | |
virtual std::string | clusterid (int timeout_ms)=0 |
Returns the ClusterId as reported in broker metadata. More... | |
virtual struct rd_kafka_s * | c_ptr ()=0 |
Returns the underlying librdkafka C rd_kafka_t handle. More... | |
virtual int32_t | controllerid (int timeout_ms)=0 |
Returns the current ControllerId (controller broker id) as reported in broker metadata. More... | |
virtual ErrorCode | fatal_error (std::string &errstr) const =0 |
Returns the first fatal error set on this client instance, or ERR_NO_ERROR if no fatal error has occurred. More... | |
virtual ErrorCode | oauthbearer_set_token (const std::string &token_value, int64_t md_lifetime_ms, const std::string &md_principal_name, const std::list< std::string > &extensions, std::string &errstr)=0 |
Set SASL/OAUTHBEARER token and metadata. More... | |
virtual ErrorCode | oauthbearer_set_token_failure (const std::string &errstr)=0 |
SASL/OAUTHBEARER token refresh failure indicator. More... | |
virtual Error * | sasl_background_callbacks_enable ()=0 |
Enable SASL OAUTHBEARER refresh callbacks on the librdkafka background thread. More... | |
virtual Queue * | get_sasl_queue ()=0 |
virtual Queue * | get_background_queue ()=0 |
virtual void * | mem_malloc (size_t size)=0 |
Allocate memory using the same allocator librdkafka uses. More... | |
virtual void | mem_free (void *ptr)=0 |
Free pointer returned by librdkafka. More... | |
virtual Error * | sasl_set_credentials (const std::string &username, const std::string &password)=0 |
Sets SASL credentials used for SASL PLAIN and SCRAM mechanisms by this Kafka client. More... | |
Base handle, super class for specific clients.
|
pure virtual |
|
pure virtual |
Returns the client's broker-assigned group member id.
|
pure virtual |
Polls the provided kafka handle for events.
Events will trigger application provided callbacks to be called.
The timeout_ms
argument specifies the maximum amount of time (in milliseconds) that the call will block waiting for events. For non-blocking calls, provide 0 as timeout_ms
. To wait indefinately for events, provide -1.
Events:
|
pure virtual |
Returns the current out queue length.
The out queue contains messages and requests waiting to be sent to, or acknowledged by, the broker.
|
pure virtual |
Request Metadata from broker.
Parameters: all_topics
- if non-zero: request info about all topics in cluster, if zero: only request info about locally known topics. only_rkt
- only request info about this topic metadatap
- pointer to hold metadata result. The *metadatap
pointer must be released with delete
. timeout_ms
- maximum response time before failing.
*metadatap
will be set), else RdKafka::ERR__TIMED_OUT on timeout or other error code on error.
|
pure virtual |
Pause producing or consumption for the provided list of partitions.
Success or error is returned per-partition in the partitions
list.
|
pure virtual |
Resume producing or consumption for the provided list of partitions.
Success or error is returned per-partition in the partitions
list.
|
pure virtual |
Query broker for low (oldest/beginning) and high (newest/end) offsets for partition.
Offsets are returned in *low
and *high
respectively.
|
pure virtual |
Get last known low (oldest/beginning) and high (newest/end) offsets for partition.
The low offset is updated periodically (if statistics.interval.ms is set) while the high offset is updated on each fetched message set from the broker.
If there is no cached offset (either low or high, or both) then OFFSET_INVALID will be returned for the respective offset.
Offsets are returned in *low
and *high
respectively.
|
pure virtual |
Look up the offsets for the given partitions by timestamp.
The returned offset for each partition is the earliest offset whose timestamp is greater than or equal to the given timestamp in the corresponding partition.
The timestamps to query are represented as offset
in offsets
on input, and offset()
will return the closest earlier offset for the timestamp on output.
Timestamps are expressed as milliseconds since epoch (UTC).
The function will block for at most timeout_ms
milliseconds.
err()
|
pure virtual |
Retrieve queue for a given partition.
Forward librdkafka logs (and debug) to the specified queue for serving with one of the ..poll() calls.
This allows an application to serve log callbacks (log_cb
) in its thread of choice.
queue | Queue to forward logs to. If the value is NULL the logs are forwarded to the main queue. |
log.queue
MUST also be set to true.
|
pure virtual |
Cancels the current callback dispatcher (Handle::poll(), KafkaConsumer::consume(), etc).
A callback may use this to force an immediate return to the calling code (caller of e.g. Handle::poll()) without processing any further events.
|
pure virtual |
Returns the ClusterId as reported in broker metadata.
timeout_ms | If there is no cached value from metadata retrieval then this specifies the maximum amount of time (in milliseconds) the call will block waiting for metadata to be retrieved. Use 0 for non-blocking calls. |
|
pure virtual |
Returns the underlying librdkafka C rd_kafka_t handle.
rd_kafka_t*
|
pure virtual |
Returns the current ControllerId (controller broker id) as reported in broker metadata.
timeout_ms | If there is no cached value from metadata retrieval then this specifies the maximum amount of time (in milliseconds) the call will block waiting for metadata to be retrieved. Use 0 for non-blocking calls. |
|
pure virtual |
Returns the first fatal error set on this client instance, or ERR_NO_ERROR if no fatal error has occurred.
This function is to be used with the Idempotent Producer and the Event class for EVENT_ERROR
events to detect fatal errors.
Generally all errors raised by the error event are to be considered informational and temporary, the client will try to recover from all errors in a graceful fashion (by retrying, etc).
However, some errors should logically be considered fatal to retain consistency; in particular a set of errors that may occur when using the Idempotent Producer and the in-order or exactly-once producer guarantees can't be satisfied.
errstr | A human readable error string if a fatal error was set. |
|
pure virtual |
Set SASL/OAUTHBEARER token and metadata.
token_value | the mandatory token value to set, often (but not necessarily) a JWS compact serialization as per https://tools.ietf.org/html/rfc7515#section-3.1. |
md_lifetime_ms | when the token expires, in terms of the number of milliseconds since the epoch. |
md_principal_name | the Kafka principal name associated with the token. |
extensions | potentially empty SASL extension keys and values where element [i] is the key and [i+1] is the key's value, to be communicated to the broker as additional key-value pairs during the initial client response as per https://tools.ietf.org/html/rfc7628#section-3.1. The number of SASL extension keys plus values must be a non-negative multiple of 2. Any provided keys and values are copied. |
errstr | A human readable error string is written here, only if there is an error. |
The SASL/OAUTHBEARER token refresh callback should invoke this method upon success. The extension keys must not include the reserved key "`auth`", and all extension keys and values must conform to the required format as per https://tools.ietf.org/html/rfc7628#section-3.1:
key = 1*(ALPHA) value = *(VCHAR / SP / HTAB / CR / LF )
RdKafka::ERR_NO_ERROR
on success, otherwise errstr
set and:RdKafka::ERR__INVALID_ARG
if any of the arguments are invalid;RdKafka::ERR__NOT_IMPLEMENTED
if SASL/OAUTHBEARER is not supported by this build;RdKafka::ERR__STATE
if SASL/OAUTHBEARER is supported but is not configured as the client's authentication mechanism."oauthbearer_token_refresh_cb"
|
pure virtual |
SASL/OAUTHBEARER token refresh failure indicator.
errstr | human readable error reason for failing to acquire a token. |
The SASL/OAUTHBEARER token refresh callback should invoke this method upon failure to refresh the token.
RdKafka::ERR_NO_ERROR
on success, otherwise:RdKafka::ERR__NOT_IMPLEMENTED
if SASL/OAUTHBEARER is not supported by this build;RdKafka::ERR__STATE
if SASL/OAUTHBEARER is supported but is not configured as the client's authentication mechanism."oauthbearer_token_refresh_cb"
|
pure virtual |
Enable SASL OAUTHBEARER refresh callbacks on the librdkafka background thread.
This serves as an alternative for applications that do not call RdKafka::Handle::poll() (et.al.) at regular intervals.
|
pure virtual |
|
pure virtual |
|
pure virtual |
Allocate memory using the same allocator librdkafka uses.
This is typically an abstraction for the malloc(3) call and makes sure the application can use the same memory allocator as librdkafka for allocating pointers that are used by librdkafka.
|
pure virtual |
Free pointer returned by librdkafka.
This is typically an abstraction for the free(3) call and makes sure the application can use the same memory allocator as librdkafka for freeing pointers returned by librdkafka.
In standard setups it is usually not necessary to use this interface rather than the free(3) function.
|
pure virtual |
Sets SASL credentials used for SASL PLAIN and SCRAM mechanisms by this Kafka client.
This function sets or resets the SASL username and password credentials used by this Kafka client. The new credentials will be used the next time this client needs to authenticate to a broker. will not disconnect existing connections that might have been made using the old credentials.