NVIDIA DeepStream SDK API Reference

9.1 Release
sources/libs/kafka_protocol_adaptor/kafka_client.h
Go to the documentation of this file.
1 /*
2  * SPDX-FileCopyrightText: Copyright (c) 2018-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
3  * SPDX-License-Identifier: Apache-2.0
4  *
5  * Licensed under the Apache License, Version 2.0 (the "License");
6  * you may not use this file except in compliance with the License.
7  * You may obtain a copy of the License at
8  *
9  * http://www.apache.org/licenses/LICENSE-2.0
10  *
11  * Unless required by applicable law or agreed to in writing, software
12  * distributed under the License is distributed on an "AS IS" BASIS,
13  * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14  * See the License for the specific language governing permissions and
15  * limitations under the License.
16  */
17 
18 #include <iostream>
19 #include <vector>
20 #include <string>
21 #include <sstream>
22 using namespace std;
23 
24 #include "rdkafka.h"
25 #include "nvds_msgapi.h"
26 
27 #define MAX_FIELD_LEN 1024
28 #define MAX_TOPIC_LEN 255 //maximum topic length supported by kafka is 255
29 #define NVDS_KAFKA_LOG_CAT "DSLOG:NVDS_KAFKA_PROTO"
30 
32  public:
33  virtual void sendcomplete(NvDsMsgApiErrorType);
34  NvDsMsgApiErrorType get_err();
35  virtual ~NvDsKafkaSendCompl() = default;
36 };
37 
39  private:
40  uint8_t *compl_flag;
42 
43  public:
44  NvDsKafkaSyncSendCompl(uint8_t *);
45  void sendcomplete(NvDsMsgApiErrorType);
46  NvDsMsgApiErrorType get_err();
47 };
48 
50  private:
51  void *user_ptr;
52  nvds_msgapi_send_cb_t async_send_cb;
53 
54  public:
56  void sendcomplete(NvDsMsgApiErrorType);
57 };
58 
59 typedef struct {
60  pthread_t consumer_tid; /* Thread which waits on incoming msg from cloud*/
62  void *user_ptr; /* User context pointer */
64 
65 typedef struct {
66  rd_kafka_t *consumer; /* Consumer instance handle */
67  char consumer_grp_id[MAX_FIELD_LEN]; /* Consumer group id */
68  consumer_thread_info cinfo; /* Consumer thread info*/
69  bool disconnect; /* variable to notify consume thread to quit */
70  string config; /* config options for consumer instance */
72 
73 typedef struct {
74  char partition_key_field[MAX_FIELD_LEN]; /* partition key for messages */
75  rd_kafka_t *producer; /* Producer instance handle */
77 
78 typedef struct {
79  char brokers[MAX_FIELD_LEN]; /* Broker string - comma separated host:port */
80  producer_instance_t p_instance; /* Producer instance details */
81  consumer_instance_t c_instance; /* consumer instance details */
83 
85 NvDsMsgApiErrorType nvds_kafka_producer_launch(void *kh, rd_kafka_conf_t *conf);
86 NvDsMsgApiErrorType nvds_kafka_client_send(void *kh, const uint8_t *payload, int len, char *topic, int sync, void *ctx, nvds_msgapi_send_cb_t cb, char *key, int keylen);
87 NvDsMsgApiErrorType nvds_kafka_client_setconf(rd_kafka_conf_t *conf, char *key, char *val);
88 void nvds_kafka_client_poll(void *kv);
89 void nvds_kafka_client_finish(void *kv);
producer_instance_t::producer
rd_kafka_t * producer
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:75
nvds_kafka_client_init
void * nvds_kafka_client_init(NvDsKafkaClientHandle *kh)
nvds_kafka_client_setconf
NvDsMsgApiErrorType nvds_kafka_client_setconf(rd_kafka_conf_t *conf, char *key, char *val)
MAX_FIELD_LEN
#define MAX_FIELD_LEN
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:27
NvDsKafkaSyncSendCompl
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:38
NvDsKafkaClientHandle::p_instance
producer_instance_t p_instance
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:80
NvDsMsgApiErrorType
NvDsMsgApiErrorType
Defines completion codes for operations in the messaging API.
Definition: sources/includes/nvds_msgapi.h:62
consumer_instance_t::consumer
rd_kafka_t * consumer
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:66
consumer_thread_info::consumer_tid
pthread_t consumer_tid
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:60
consumer_thread_info::subscribe_req_cb
nvds_msgapi_subscribe_request_cb_t subscribe_req_cb
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:61
nvds_kafka_client_poll
void nvds_kafka_client_poll(void *kv)
nvds_kafka_client_finish
void nvds_kafka_client_finish(void *kv)
producer_instance_t
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:73
consumer_instance_t::config
string config
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:70
nvds_kafka_producer_launch
NvDsMsgApiErrorType nvds_kafka_producer_launch(void *kh, rd_kafka_conf_t *conf)
consumer_thread_info
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:59
nvds_kafka_client_send
NvDsMsgApiErrorType nvds_kafka_client_send(void *kh, const uint8_t *payload, int len, char *topic, int sync, void *ctx, nvds_msgapi_send_cb_t cb, char *key, int keylen)
NvDsKafkaClientHandle::c_instance
consumer_instance_t c_instance
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:81
NvDsKafkaSendCompl
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:31
consumer_instance_t::cinfo
consumer_thread_info cinfo
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:68
NvDsKafkaClientHandle
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:78
nvds_msgapi_send_cb_t
void(* nvds_msgapi_send_cb_t)(void *user_ptr, NvDsMsgApiErrorType completion_flag)
Type definition for a "send" callback.
Definition: sources/includes/nvds_msgapi.h:76
consumer_instance_t
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:65
consumer_thread_info::user_ptr
void * user_ptr
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:62
nvds_msgapi_subscribe_request_cb_t
void(* nvds_msgapi_subscribe_request_cb_t)(NvDsMsgApiErrorType flag, void *msg, int msg_len, char *topic, void *user_ptr)
Type definition for callback registered during subscribe.
Definition: sources/includes/nvds_msgapi.h:92
consumer_instance_t::disconnect
bool disconnect
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:69
NvDsKafkaAsyncSendCompl
Definition: sources/libs/kafka_protocol_adaptor/kafka_client.h:49