Using Producer API to Produce Messages for Security Topics
Function
It is used by the Producer APIs to consume messages in the security topic.
Sample Code
The following is a code snippet used by the Producer API to produce messages to a secure topic.
The following code snippet belongs to the run method in the com.huawei.bigdata.kafka.example.Producer class.
/**
* The producer thread executes a function to send messages periodically.
*/
public void run() {
LOG.info("New Producer: start.");
int messageNo = 1;
while (messageNo <= MESSAGE_NUM) {
String messageStr = "Message_" + messageNo;
long startTime = System.currentTimeMillis();
// Construct a message record.
ProducerRecord<Integer, String> record = new ProducerRecord<Integer, String>(topic, messageNo, messageStr);
if (isAsync) {
// Send data asynchronously.
producer.send(record, new DemoCallBack(startTime, messageNo, messageStr));
} else {
try {
// Send data synchronously.
producer.send(record).get();
long elapsedTime = System.currentTimeMillis() - startTime;
LOG.info("message(" + messageNo + ", " + messageStr + ") sent to topic(" + topic + ") in " + elapsedTime + " ms.");
} catch (InterruptedException ie) {
LOG.info("The InterruptedException occured : {}.", ie);
} catch (ExecutionException ee) {
LOG.info("The ExecutionException occured : {}.", ee);
}
}
messageNo++;
}
} Feedback
Was this page helpful?
Provide feedbackThank you very much for your feedback. We will continue working to improve the documentation.See the reply and handling status in My Cloud VOC.
For any further questions, feel free to contact us through the chatbot.
Chatbot