- import org.apache.kafka.clients.producer.KafkaProducer;
- import org.apache.kafka.clients.producer.ProducerConfig;
- import org.apache.kafka.clients.producer.ProducerRecord;
- import org.apache.kafka.clients.producer.RecordMetadata;
- import org.apache.kafka.clients.producer.Callback;
- import org.apache.kafka.clients.producer.KafkaException;
- import org.apache.kafka.common.serialization.StringSerializer;
- import java.util.Properties;
- public class KafkaTransactionalProducer {
- public static void main(String[] args) {
- Properties props = new Properties();
- props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "10.1.1.1:9000,4.4.4.4:8888");
- props.put("security.protocol", "SASL_PLAINTEXT");
- props.put("sasl.mechanism", "SCRAM-SHA-256");
- props.put("sasl.kerberos.service.name", "target");
- props.put("sasl.jaas.config", "org.apache.kafka.common.security.scram.ScramLoginModule required username="ee" password="eres";");
- props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
- props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
- props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
- props.put(ProducerConfig.ACKS_CONFIG, "all");
- props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "my-unique-transactional-id"); // 自定义唯一事务ID
- KafkaProducer<String, String> producer = new KafkaProducer<>(props);
- try {
- producer.initTransactions(); // 初始化事务
- producer.beginTransaction(); // 开始事务
- String topic = "your_topic";
- String message = "{"key": "value"}"; // 要发送的消息
- ProducerRecord<String, String> record = new ProducerRecord<>(topic, message);
-
- producer.send(record, new Callback() {
- public void onCompletion(RecordMetadata metadata, Exception exception) {
- if (exception != null) {
- System.err.println("Error sending message: " + exception.getMessage());
- } else {
- System.out.println("Message sent successfully to topic " + metadata.topic() +
- " partition " + metadata.partition() +
- " with offset " + metadata.offset());
- }
- }
- });
- producer.commitTransaction(); // 提交事务
- } catch (KafkaException e) {
- producer.abortTransaction(); // 事务失败,回滚
- e.printStackTrace();
- } finally {
- producer.close();
- }
- }
- }
- 把这些代码改成一个独立的应用服务有自己的端口通过post请求传递json报文,然后给kafka发送消息
复制代码 要将 Kafka 生产者代码改成一个独立的应用服务,可以或许通过 POST 请求吸收 JSON 报文并将其发送到 Kafka,你可以使用 Spring Boot 框架来实现这个功能。Spring Boot 提供了便捷的 REST API 开发和设置管理功能。
以下是一个基于 Spring Boot 的示例项目,展示了怎样实现这个需求:
1. 创建 Spring Boot 项目
首先,你必要创建一个 Spring Boot 项目。如果你使用 Spring Initializr,可以选择以下依赖:
2. 添加 Maven 依赖
在 pom.xml 文件中,确保包含以下依赖
- <dependencies>
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-web</artifactId>
- </dependency>
- <dependency>
- <groupId>org.springframework.kafka</groupId>
- <artifactId>spring-kafka</artifactId>
- </dependency>
- </dependencies>
复制代码 3. 设置 Kafka 生产者
在 application.properties 文件中设置 Kafka 生产者属性:
- spring.kafka.producer.bootstrap-servers=10.1.1.1:9000,4.4.4.4:8888
- spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
- spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer
- spring.kafka.producer.enable-idempotence=true
- spring.kafka.producer.acks=all
- spring.kafka.producer.transaction-id-prefix=my-unique-transactional-id
- spring.kafka.producer.properties.security.protocol=SASL_PLAINTEXT
- spring.kafka.producer.properties.sasl.mechanism=SCRAM-SHA-256
- spring.kafka.producer.properties.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="ee" password="eres";
- spring.kafka.producer.properties.sasl.kerberos.service.name=target
复制代码 4. 实现 Kafka 生产者服务
创建一个 KafkaProducerService 类来处理惩罚消息发送:
- package com.example.kafkaproducer.service;
- import org.apache.kafka.clients.producer.KafkaProducer;
- import org.apache.kafka.clients.producer.ProducerRecord;
- import org.apache.kafka.clients.producer.RecordMetadata;
- import org.apache.kafka.clients.producer.Callback;
- import org.apache.kafka.clients.producer.KafkaException;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.kafka.core.ProducerFactory;
- import org.springframework.kafka.core.KafkaTemplate;
- import org.springframework.kafka.core.ProducerFactory;
- import org.springframework.stereotype.Service;
- @Service
- public class KafkaProducerService {
- private final KafkaTemplate<String, String> kafkaTemplate;
- @Autowired
- public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) {
- this.kafkaTemplate = kafkaTemplate;
- }
- public void sendMessage(String topic, String message) {
- kafkaTemplate.executeInTransaction(kafkaTemplate -> {
- kafkaTemplate.send(topic, message).addCallback(new org.springframework.util.concurrent.ListenableFutureCallback<>() {
- @Override
- public void onSuccess(RecordMetadata result) {
- System.out.println("Message sent successfully to topic " + result.topic() +
- " partition " + result.partition() +
- " with offset " + result.offset());
- }
- @Override
- public void onFailure(Throwable ex) {
- System.err.println("Error sending message: " + ex.getMessage());
- }
- });
- return null;
- });
- }
- }
复制代码 5. 创建 REST Controller
实现一个 REST 控制器来处理惩罚 POST 请求:
- package com.example.kafkaproducer.controller;
- import com.example.kafkaproducer.service.KafkaProducerService;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.web.bind.annotation.*;
- @RestController
- @RequestMapping("/api")
- public class KafkaController {
- private final KafkaProducerService kafkaProducerService;
- @Autowired
- public KafkaController(KafkaProducerService kafkaProducerService) {
- this.kafkaProducerService = kafkaProducerService;
- }
- @PostMapping("/send")
- public String sendMessage(@RequestBody String message) {
- String topic = "your_topic";
- kafkaProducerService.sendMessage(topic, message);
- return "Message sent to Kafka";
- }
- }
复制代码 6. 启动 Spring Boot 应用
在你的 Application 类中启动 Spring Boot 应用:
- package com.example.kafkaproducer;
- import org.springframework.boot.SpringApplication;
- import org.springframework.boot.autoconfigure.SpringBootApplication;
- @SpringBootApplication
- public class KafkaProducerApplication {
- public static void main(String[] args) {
- SpringApplication.run(KafkaProducerApplication.class, args);
- }
- }
复制代码 7. 运行应用并测试
启动 Spring Boot 应用后,你可以通过 POST 请求将 JSON 消息发送到 http://localhost:8080/api/send,好比使用 curl:
- curl -X POST http://localhost:8080/api/send -H "Content-Type: application/json" -d '{"key": "value"}'
复制代码 这样,Kafka 生产者服务就可以通过 REST API 吸收消息并将其发送到 Kafka 了。
上面的用java的1.8实现
如果你使用 Java 1.8,你可以按照以下步骤来实现一个独立的应用服务,该服务通过 POST 请求吸收 JSON 报文并将其发送到 Kafka。我们将不使用 Spring Boot,而是使用更底子的 Java Servlet 技能来完成这个使命。
1. 创建 Maven 项目
确保你的 pom.xml 包含必要的依赖:
- <dependencies>
- <!-- Spring Boot Starter Web -->
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-web</artifactId>
- </dependency>
-
- <!-- Spring Boot Starter Kafka -->
- <dependency>
- <groupId>org.springframework.kafka</groupId>
- <artifactId>spring-kafka</artifactId>
- </dependency>
- <!-- JSON处理 -->
- <dependency>
- <groupId>com.google.code.gson</groupId>
- <artifactId>gson</artifactId>
- <version>2.10.1</version>
- </dependency>
- </dependencies>
复制代码 2. 设置 Kafka 生产者
创建 KafkaProducerConfig 类来设置 Kafka 生产者:
- import org.apache.kafka.clients.producer.ProducerConfig;
- import org.apache.kafka.common.serialization.StringSerializer;
- import java.util.Properties;
- public class KafkaProducerConfig {
- public static Properties getProducerProperties() {
- Properties props = new Properties();
- props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "10.1.1.1:9000,4.4.4.4:8888");
- props.put("security.protocol", "SASL_PLAINTEXT");
- props.put("sasl.mechanism", "SCRAM-SHA-256");
- props.put("sasl.kerberos.service.name", "target");
- props.put("sasl.jaas.config", "org.apache.kafka.common.security.scram.ScramLoginModule required username="ee" password="eres";");
- props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
- props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
- props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
- props.put(ProducerConfig.ACKS_CONFIG, "all");
- props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "my-unique-transactional-id");
- return props;
- }
- }
复制代码 3. 实现 Kafka 生产者 Servlet
创建 KafkaProducerServlet 类来处理惩罚 POST 请求并将消息发送到 Kafka:
- import com.google.gson.Gson;
- import org.apache.kafka.clients.producer.KafkaProducer;
- import org.apache.kafka.clients.producer.ProducerRecord;
- import org.apache.kafka.clients.producer.RecordMetadata;
- import org.apache.kafka.clients.producer.Callback;
- import org.apache.kafka.clients.producer.KafkaException;
- import javax.servlet.ServletException;
- import javax.servlet.annotation.WebServlet;
- import javax.servlet.http.HttpServlet;
- import javax.servlet.http.HttpServletRequest;
- import javax.servlet.http.HttpServletResponse;
- import java.io.IOException;
- @WebServlet("/api/send")
- public class KafkaProducerServlet extends HttpServlet {
- private KafkaProducer<String, String> producer;
- @Override
- public void init() throws ServletException {
- super.init();
- producer = new KafkaProducer<>(KafkaProducerConfig.getProducerProperties());
- }
- @Override
- protected void doPost(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {
- String topic = "your_topic";
- StringBuilder stringBuilder = new StringBuilder();
- String line;
- while ((line = req.getReader().readLine()) != null) {
- stringBuilder.append(line);
- }
- String message = stringBuilder.toString();
- ProducerRecord<String, String> record = new ProducerRecord<>(topic, message);
- producer.send(record, new Callback() {
- @Override
- public void onCompletion(RecordMetadata metadata, Exception exception) {
- if (exception != null) {
- System.err.println("Error sending message: " + exception.getMessage());
- } else {
- System.out.println("Message sent successfully to topic " + metadata.topic() +
- " partition " + metadata.partition() +
- " with offset " + metadata.offset());
- }
- }
- });
- resp.setContentType("text/plain");
- resp.getWriter().write("Message sent to Kafka");
- }
- @Override
- public void destroy() {
- super.destroy();
- if (producer != null) {
- producer.close();
- }
- }
- }
复制代码 4. 设置 web.xml
如果你不使用注解,可以在 web.xml 中设置 Servlet:
- <web-app xmlns="http://xmlns.jcp.org/xml/ns/javaee"
- xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
- xsi:schemaLocation="http://xmlns.jcp.org/xml/ns/javaee
- http://xmlns.jcp.org/xml/ns/javaee/web-app_4_0.xsd"
- version="4.0">
- <servlet>
- <servlet-name>KafkaProducerServlet</servlet-name>
- <servlet-class>com.example.KafkaProducerServlet</servlet-class>
- </servlet>
- <servlet-mapping>
- <servlet-name>KafkaProducerServlet</servlet-name>
- <url-pattern>/api/send</url-pattern>
- </servlet-mapping>
- </web-app>
复制代码 5. 部署和运行
将你的应用打包成一个 WAR 文件并部署到支持 Servlet 的服务器(如 Apache Tomcat)。然后,你可以通过 POST 请求将 JSON 消息发送到 http://localhost:8080/your-app/api/send。
例如,使用 curl 来测试:
- bashcurl -X POST http://localhost:8080/your-app/api/send -H "Content-Type: application/json" -d '{"key": "value"}'
复制代码 这样,你就完成了将 Kafka 生产者代码改成一个独立的应用服务,通过 POST 请求吸收 JSON 报文并发送到 Kafka 的使命。
上面的用java的1.8springboot实现
要使用 Java 1.8 和 Spring Boot 实现一个独立的应用服务,通过 POST 请求吸收 JSON 报文并将其发送到 Kafka,你可以按照以下步骤举行。
1. 创建 Spring Boot 项目
可以使用 Spring Initializr 来创建一个新的 Spring Boot 项目。选择以下依赖:
2. 添加 Maven 依赖
在 pom.xml 文件中添加 Kafka 和 Spring Boot 相干的依赖:
- <dependencies>
- <!-- Spring Boot Starter Web -->
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-web</artifactId>
- </dependency>
-
- <!-- Spring Boot Starter Kafka -->
- <dependency>
- <groupId>org.springframework.kafka</groupId>
- <artifactId>spring-kafka</artifactId>
- </dependency>
- <!-- JSON处理 -->
- <dependency>
- <groupId>com.google.code.gson</groupId>
- <artifactId>gson</artifactId>
- <version>2.10.1</version>
- </dependency>
- </dependencies>
复制代码 3. 设置 Kafka 生产者
在 application.properties 文件中设置 Kafka 生产者属性:
- spring.kafka.producer.bootstrap-servers=10.1.1.1:9000,4.4.4.4:8888
- spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
- spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer
- spring.kafka.producer.enable-idempotence=true
- spring.kafka.producer.acks=all
- spring.kafka.producer.transaction-id-prefix=my-unique-transactional-id
- spring.kafka.producer.properties.security.protocol=SASL_PLAINTEXT
- spring.kafka.producer.properties.sasl.mechanism=SCRAM-SHA-256
- spring.kafka.producer.properties.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="ee" password="eres";
- spring.kafka.producer.properties.sasl.kerberos.service.name=target
复制代码 4. 创建 Kafka 生产者服务
创建一个 KafkaProducerService 类来处理惩罚消息发送:
- package com.example.kafkaproducer.service;
- import org.springframework.kafka.core.KafkaTemplate;
- import org.springframework.kafka.core.ProducerFactory;
- import org.springframework.kafka.core.KafkaTemplate;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.kafka.annotation.KafkaListener;
- import org.springframework.kafka.core.KafkaTemplate;
- import org.springframework.stereotype.Service;
- @Service
- public class KafkaProducerService {
- private final KafkaTemplate<String, String> kafkaTemplate;
- @Autowired
- public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) {
- this.kafkaTemplate = kafkaTemplate;
- }
- public void sendMessage(String topic, String message) {
- kafkaTemplate.executeInTransaction(kafkaTemplate -> {
- kafkaTemplate.send(topic, message).addCallback(new ListenableFutureCallback<>() {
- @Override
- public void onSuccess(RecordMetadata result) {
- System.out.println("Message sent successfully to topic " + result.topic() +
- " partition " + result.partition() +
- " with offset " + result.offset());
- }
- @Override
- public void onFailure(Throwable ex) {
- System.err.println("Error sending message: " + ex.getMessage());
- }
- });
- return null;
- });
- }
- }
复制代码 5. 创建 REST Controller
创建一个 REST 控制器来处理惩罚 POST 请求:
- package com.example.kafkaproducer.controller;
- import com.example.kafkaproducer.service.KafkaProducerService;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.web.bind.annotation.*;
- @RestController
- @RequestMapping("/api")
- public class KafkaController {
- private final KafkaProducerService kafkaProducerService;
- @Autowired
- public KafkaController(KafkaProducerService kafkaProducerService) {
- this.kafkaProducerService = kafkaProducerService;
- }
- @PostMapping("/send")
- public String sendMessage(@RequestBody String message) {
- String topic = "your_topic";
- kafkaProducerService.sendMessage(topic, message);
- return "Message sent to Kafka";
- }
- }
复制代码 6. 启动 Spring Boot 应用
在你的 Application 类中启动 Spring Boot 应用:
- package com.example.kafkaproducer;
- import org.springframework.boot.SpringApplication;
- import org.springframework.boot.autoconfigure.SpringBootApplication;
- @SpringBootApplication
- public class KafkaProducerApplication {
- public static void main(String[] args) {
- SpringApplication.run(KafkaProducerApplication.class, args);
- }
- }
复制代码 7. 运行应用并测试
启动 Spring Boot 应用后,你可以通过 POST 请求将 JSON 消息发送到 http://localhost:8080/api/send,例如使用 curl:
- package com.example.kafkaproducer;
- import org.springframework.boot.SpringApplication;
- import org.springframework.boot.autoconfigure.SpringBootApplication;
- @SpringBootApplication
- public class KafkaProducerApplication {
- public static void main(String[] args) {
- SpringApplication.run(KafkaProducerApplication.class, args);
- }
- }
复制代码 额外提示
- 非常处理惩罚: 在实际应用中,你大概必要添加更完善的非常处理惩罚和日志记录。
- 安全性: 确保你的应用具有恰当的安全设置,尤其是对外暴露的 API。
这样,你就完成了使用 Java 1.8 和 Spring Boot 实现 Kafka 生产者服务的使命,通过 POST 请求吸收 JSON 报文并将其发送到 Kafka。
public void onSuccess(RecordMetadata result) { 我这个继承的参数不是RecordMetadata result
如果你的 onSuccess 方法的参数不是 RecordMetadata result,那么你大概在使用 ListenableFutureCallback 接口的实现时遇到了标题。在 Kafka 的生产者 API 中,ListenableFutureCallback 接口通常用于处理惩罚 send 方法的回调。
请确保你正确地实现了 ListenableFutureCallback 接口。以下是一个完整的示例,展示了怎样实现 Kafka 生产者的回调处理惩罚,包括 onSuccess 和 onFailure 方法。
完整示例:KafkaProducerService 类
在这个示例中,我们将使用 ListenableFutureCallback 处理惩罚 Kafka 生产的结果。
- package com.example.kafkaproducer.service;
- import org.apache.kafka.clients.producer.KafkaProducer;
- import org.apache.kafka.clients.producer.ProducerRecord;
- import org.apache.kafka.clients.producer.RecordMetadata;
- import org.apache.kafka.clients.producer.Callback;
- import org.apache.kafka.clients.producer.KafkaException;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.kafka.core.KafkaTemplate;
- import org.springframework.kafka.core.ProducerFactory;
- import org.springframework.kafka.core.KafkaTemplate;
- import org.springframework.stereotype.Service;
- import org.springframework.util.concurrent.ListenableFutureCallback;
- @Service
- public class KafkaProducerService {
- private final KafkaTemplate<String, String> kafkaTemplate;
- @Autowired
- public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) {
- this.kafkaTemplate = kafkaTemplate;
- }
- public void sendMessage(String topic, String message) {
- kafkaTemplate.executeInTransaction(kafkaTemplate -> {
- kafkaTemplate.send(topic, message).addCallback(new ListenableFutureCallback<>() {
- @Override
- public void onSuccess(RecordMetadata result) {
- System.out.println("Message sent successfully to topic " + result.topic() +
- " partition " + result.partition() +
- " with offset " + result.offset());
- }
- @Override
- public void onFailure(Throwable ex) {
- System.err.println("Error sending message: " + ex.getMessage());
- }
- });
- return null;
- });
- }
- }
复制代码 关键点
- KafkaTemplate 的回调: 使用 kafkaTemplate.send().addCallback 方法可以处理惩罚消息发送的结果。addCallback 方法接受一个 ListenableFutureCallback 对象来处理惩罚乐成和失败的回调。
- 实现 ListenableFutureCallback: ListenableFutureCallback 是 Spring 提供的接口,用于处理惩罚异步操作的乐成和失败结果。你必要实现 onSuccess 和 onFailure 方法。
- ProducerRecord 和 RecordMetadata: ProducerRecord 用于构建消息记录,RecordMetadata 用于吸收发送消息后的元数据。
使用 Callback 接口的示例
如果你直接使用 Kafka 的 Callback 接口而不是 Spring 的 ListenableFutureCallback,示例如下:
- import org.apache.kafka.clients.producer.Callback;
- import org.apache.kafka.clients.producer.RecordMetadata;
- public class KafkaProducerCallback implements Callback {
- @Override
- public void onCompletion(RecordMetadata metadata, Exception exception) {
- if (exception != null) {
- System.err.println("Error sending message: " + exception.getMessage());
- } else {
- System.out.println("Message sent successfully to topic " + metadata.topic() +
- " partition " + metadata.partition() +
- " with offset " + metadata.offset());
- }
- }
- }
复制代码 然后在 KafkaProducerService 中使用这个 Callback 类:
- public void sendMessage(String topic, String message) {
- kafkaTemplate.executeInTransaction(kafkaTemplate -> {
- kafkaTemplate.send(topic, message).addCallback(new KafkaProducerCallback());
- return null;
- });
- }
复制代码 这样,你就可以确保 Kafka 消息发送的结果被正确地处理惩罚,无论是乐成还是失败。
public void onSuccess(RecordMetadata result) { 我这个继承的参数shi Object o怎么回事
如果你在实现 Kafka 生产者的回调时遇到了 onSuccess 方法的参数类型为 Object o 的标题,这大概是由于使用了差别的回调接口或设置。
在 Kafka 中,send 方法的回调接口通常是 Callback,其 onCompletion 方法的签名如下:
void onCompletion(RecordMetadata metadata, Exception exception);
如果你的实现中 onSuccess 方法的参数类型是 Object,大概你在使用的是与 Kafka 原生 API 不匹配的回调接口或在自定义实现时出现了标题。
确保使用正确的回调接口
以下是怎样正确实现 Kafka 的 Callback 接口,并处理惩罚 onSuccess 和 onFailure 方法:
1. Kafka Callback 接口的正确实现
- import org.apache.kafka.clients.producer.Callback;
- import org.apache.kafka.clients.producer.RecordMetadata;
- import org.apache.kafka.clients.producer.ProducerRecord;
- public class KafkaProducerCallback implements Callback {
- @Override
- public void onCompletion(RecordMetadata metadata, Exception exception) {
- if (exception != null) {
- System.err.println("Error sending message: " + exception.getMessage());
- } else {
- System.out.println("Message sent successfully to topic " + metadata.topic() +
- " partition " + metadata.partition() +
- " with offset " + metadata.offset());
- }
- }
- }
复制代码 2. 使用 Kafka Callback 的 send 方法
在 Kafka 生产者服务中使用上述回调:
- import org.apache.kafka.clients.producer.KafkaProducer;
- import org.apache.kafka.clients.producer.ProducerRecord;
- import org.apache.kafka.clients.producer.ProducerConfig;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.kafka.core.KafkaTemplate;
- import org.springframework.stereotype.Service;
- @Service
- public class KafkaProducerService {
- private final KafkaTemplate<String, String> kafkaTemplate;
- @Autowired
- public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) {
- this.kafkaTemplate = kafkaTemplate;
- }
- public void sendMessage(String topic, String message) {
- kafkaTemplate.executeInTransaction(kafkaTemplate -> {
- kafkaTemplate.send(topic, message).addCallback(new KafkaProducerCallback());
- return null;
- });
- }
- }
复制代码 自定义实现中的 Object o 标题
如果你看到 onSuccess(Object o),这大概是你在自定义回调或在其他上下文中使用了差别的回调接口。请查抄以下几点:
- 查抄回调接口类型:确保你的回调接口实现的是 Kafka 的 Callback 接口,而不是其他的接口或自定义接口。
- 回调方法的签名:确保 onSuccess 方法的签名与 Callback 接口的要求匹配。如果你的代码中 onSuccess 方法接受 Object 类型的参数,大概必要更改为 RecordMetadata 类型。
- 代码示例:确保在代码中,send 方法的回调处理惩罚与 Kafka 的回调接口划一。
结论
确保你使用的回调接口正确且实现了 Kafka 的 Callback 接口,以便 onCompletion 方法可以正确处理惩罚消息的发送结果。如果你遇到差别类型的回调接口,大概必要查抄代码的其他部分或依赖项以确保划一
免责声明:如果侵犯了您的权益,请联系站长,我们会及时删除侵权内容,谢谢合作!更多信息从访问主页:qidao123.com:ToB企服之家,中国第一个企服评测及商务社交产业平台。 |