From e4915995f491d7ad32cdc6cac929cf8ea7be4a28 Mon Sep 17 00:00:00 2001 From: YunaiV <> Date: Mon, 9 Mar 2020 22:39:56 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=20spring=20cloud=20stream=20?= =?UTF-8?q?kafka=20=E7=A4=BA=E4=BE=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/main/resources/application.yml | 2 +- .../kafkademo/config/TransactionConfig.java | 10 +++------- .../kafkademo/controller/Demo01Controller.java | 12 ------------ .../src/main/resources/application.yml | 6 +++--- 4 files changed, 7 insertions(+), 23 deletions(-) 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: