|
1 | 1 | /* |
2 | | - * Copyright 2016-2020 the original author or authors. |
| 2 | + * Copyright 2016-2021 the original author or authors. |
3 | 3 | * |
4 | 4 | * Licensed under the Apache License, Version 2.0 (the "License"); |
5 | 5 | * you may not use this file except in compliance with the License. |
|
58 | 58 | import org.springframework.context.ApplicationListener; |
59 | 59 | import org.springframework.context.event.ContextStoppedEvent; |
60 | 60 | import org.springframework.core.log.LogAccessor; |
| 61 | +import org.springframework.kafka.KafkaException; |
61 | 62 | import org.springframework.kafka.support.TransactionSupport; |
62 | 63 | import org.springframework.lang.Nullable; |
63 | 64 | import org.springframework.util.Assert; |
@@ -716,7 +717,20 @@ private CloseSafeProducer<K, V> doCreateTxProducer(String prefix, String suffix, |
716 | 717 | } |
717 | 718 | checkBootstrap(newProducerConfigs); |
718 | 719 | newProducer = createRawProducer(newProducerConfigs); |
719 | | - newProducer.initTransactions(); |
| 720 | + try { |
| 721 | + newProducer.initTransactions(); |
| 722 | + } |
| 723 | + catch (RuntimeException ex) { |
| 724 | + try { |
| 725 | + newProducer.close(this.physicalCloseTimeout); |
| 726 | + } |
| 727 | + catch (RuntimeException ex2) { |
| 728 | + KafkaException newEx = new KafkaException("initTransactions() failed and then close() failed", ex); |
| 729 | + newEx.addSuppressed(ex2); |
| 730 | + throw newEx; // NOSONAR - lost stack trace |
| 731 | + } |
| 732 | + throw new KafkaException("initTransactions() failed", ex); |
| 733 | + } |
720 | 734 | CloseSafeProducer<K, V> closeSafeProducer = |
721 | 735 | new CloseSafeProducer<>(newProducer, remover, prefix, this.physicalCloseTimeout, this.beanName, |
722 | 736 | this.epoch.get()); |
|
0 commit comments