diff --git a/labx-11/labx-11-sc-stream-kafka-consumer-transaction/src/main/resources/application.yml b/labx-11/labx-11-sc-stream-kafka-consumer-transaction/src/main/resources/application.yml index e8050e7a..ce42cb94 100644 --- a/labx-11/labx-11-sc-stream-kafka-consumer-transaction/src/main/resources/application.yml +++ b/labx-11/labx-11-sc-stream-kafka-consumer-transaction/src/main/resources/application.yml @@ -25,7 +25,7 @@ spring: consumer: configuration: isolation: - level: read_committed + level: read_committed # 读取已提交的消息 server: port: ${random.int[10000,19999]} # 随机端口,方便启动多个消费者 diff --git a/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/config/TransactionConfig.java b/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/config/TransactionConfig.java index d9505ba7..9acaf867 100644 --- a/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/config/TransactionConfig.java +++ b/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/config/TransactionConfig.java @@ -14,18 +14,14 @@ import org.springframework.transaction.annotation.EnableTransactionManagement; @EnableTransactionManagement public class TransactionConfig { -// @Bean -// public KafkaTransactionManager transactionManager(ProducerFactory producerFactory) { -// return new KafkaTransactionManager<>(producerFactory); -// } - @Bean public PlatformTransactionManager transactionManager(BinderFactory binders) { + // 获得 Kafka ProducerFactory 对象 ProducerFactory pf = ((KafkaMessageChannelBinder) binders.getBinder(null, MessageChannel.class)).getTransactionalProducerFactory(); + // 创建 KafkaTransactionManager 事务管理器 + assert pf != null; return new KafkaTransactionManager<>(pf); } - - } diff --git a/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java b/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java index 55c06c86..0faa1581 100644 --- a/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java +++ b/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/java/cn/iocoder/springcloud/labx11/kafkademo/kafkademo/controller/Demo01Controller.java @@ -23,18 +23,6 @@ public class Demo01Controller { @Autowired private MySource mySource; - @GetMapping("/send") - public boolean send() { - // 创建 Message - Demo01Message message = new Demo01Message() - .setId(new Random().nextInt()); - // 创建 Spring Message 对象 - Message springMessage = MessageBuilder.withPayload(message) - .build(); - // 发送消息 - return mySource.demo01Output().send(springMessage); - } - @Transactional @GetMapping("/send_transaction") public void sendTransaction() throws InterruptedException { diff --git a/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/resources/application.yml b/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/resources/application.yml index 8a5a1ca2..94e7061c 100644 --- a/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/resources/application.yml +++ b/labx-11/labx-11-sc-stream-kafka-producer-transaction/src/main/resources/application.yml @@ -17,11 +17,11 @@ spring: binder: brokers: 127.0.0.1:9092 # 指定 Kafka Broker 地址,可以设置多个,以逗号分隔 transaction: - transaction-id-prefix: demo. # TODO + transaction-id-prefix: demo. # 事务编号前缀 producer: configuration: - retries: 1 - acks: all + retries: 1 # 发送失败时,重试发送的次数 + acks: all # 0-不应答。1-leader 应答。all-所有 leader 和 follower 应答。 # Kafka 自定义 Binding 配置项,对应 KafkaBindingProperties Map bindings: demo01-output: