ChainedTransactionManager with DataSourceTransactionManager + KafkaTransactionManager. Kafka sends operation are not getting rolling back

Viewed 322

I am trying perform DB + kafka send operation in @Transactional. But when exception occur in transactional method, DB operation is getting rollback successfully but kafka transaction is not getting rollback.

I am using spring-kafka - 2.5.2.RELEASE. According to this thread https://github.com/spring-projects/spring-kafka/issues/433, JpaTransactionManager + KafkaTransactionManager , is working

    @EnableTransactionManagement
    public class SampleBaseApplication {
    
        public static void main(String[] args) {
            ConfigurableApplicationContext ctx = SpringApplication.run(SampleBaseApplication.class, args);
            Map<String, PlatformTransactionManager> tms = ctx.getBeansOfType(PlatformTransactionManager.class);
            System.out.println(tms);
            ctx.close();
        }
    
        @Bean
        public ApplicationRunner runner1(Foo foo) {
            return args -> foo.sendToKafkaAndDB();
        }
    
        @Bean
        public DataSourceTransactionManager dstm(DataSource dataSource) {
            return new DataSourceTransactionManager(dataSource);
        }
    
        @Bean(name="chainedTxMang")
        public ChainedTransactionManager chainedTxM(DataSourceTransactionManager dstm, KafkaTransactionManager<?, ?> kafka) {
            return new ChainedTransactionManager(dstm, kafka);
        }
    
        @Component
        public static class Foo {
    
            @Autowired
            @Qualifier("transactionalTemplate")
            private KafkaTemplate<String, String> template;
            
            @Autowired
            private DataAccess dataAccess;
            
            @Transactional(transactionManager = "chainedTxMang")
            public void sendToKafkaAndDB() throws Exception {
                dataAccess.insertInTable("113", "111",
                        "111", "COMPLETED", "113");
                
                System.out.println(this.template.send("TEST_TOPIC", "bar").get());
                
                throw new RuntimeException("exp...");
            }
    
        }
    
    }

below are the logs -

USER\Downloads\sample\target\classes started by DC-USER in C:\Users\DC-USER\Downloads\sample)
2021-Jun-21 20:17:44 PM [main] INFO  com.example.SampleBaseApplication -  {} - No active profile set, falling back to default profiles: default
2021-Jun-21 20:17:49 PM [main] INFO  org.springframework.boot.web.embedded.tomcat.TomcatWebServer -  {} - Tomcat initialized with port(s): 8082 (http)
2021-Jun-21 20:17:49 PM [main] INFO  org.apache.coyote.http11.Http11NioProtocol -  {} - Initializing ProtocolHandler ["http-nio-8082"]
2021-Jun-21 20:17:49 PM [main] INFO  org.apache.catalina.core.StandardService -  {} - Starting service [Tomcat]
2021-Jun-21 20:17:49 PM [main] INFO  org.apache.catalina.core.StandardEngine -  {} - Starting Servlet engine: [Apache Tomcat/9.0.36]
2021-Jun-21 20:17:53 PM [main] INFO  org.apache.catalina.core.ContainerBase.[Tomcat].[localhost].[/] -  {} - Initializing Spring embedded WebApplicationContext
2021-Jun-21 20:17:53 PM [main] INFO  org.springframework.boot.web.servlet.context.ServletWebServerApplicationContext -  {} - Root WebApplicationContext: initialization completed in 8303 ms
2021-Jun-21 20:17:54 PM [main] INFO  org.springframework.aop.framework.CglibAopProxy -  {} - Unable to proxy interface-implementing method [public final void org.springframework.dao.support.DaoSupport.afterPropertiesSet() throws java.lang.IllegalArgumentException,org.springframework.beans.factory.BeanInitializationException] because it is marked as final: Consider using interface-based JDK proxies instead!
2021-Jun-21 20:17:54 PM [main] TRACE org.springframework.transaction.annotation.AnnotationTransactionAttributeSource -  {} - Adding transactional method 'com.example.SampleBaseApplication$Foo.sendToKafkaAndDB' with attribute: PROPAGATION_REQUIRED,ISOLATION_DEFAULT; 'chainedTxMang'
2021-Jun-21 20:17:55 PM [main] INFO  org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor -  {} - Initializing ExecutorService 'applicationTaskExecutor'
2021-Jun-21 20:17:58 PM [main] INFO  org.springframework.boot.actuate.endpoint.web.EndpointLinksResolver -  {} - Exposing 14 endpoint(s) beneath base path '/actuator'
2021-Jun-21 20:17:58 PM [main] INFO  org.apache.coyote.http11.Http11NioProtocol -  {} - Starting ProtocolHandler ["http-nio-8082"]
2021-Jun-21 20:17:58 PM [main] INFO  org.springframework.boot.web.embedded.tomcat.TomcatWebServer -  {} - Tomcat started on port(s): 8082 (http) with context path ''
2021-Jun-21 20:17:58 PM [main] INFO  com.example.SampleBaseApplication -  {} - Started SampleBaseApplication in 15.457 seconds (JVM running for 22.468)
2021-Jun-21 20:17:58 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Initializing transaction synchronization
2021-Jun-21 20:17:58 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Clearing transaction synchronization
2021-Jun-21 20:17:58 PM [main] INFO  com.zaxxer.hikari.HikariDataSource -  {} - Hikari Handler DB Pool - Starting...
2021-Jun-21 20:17:59 PM [main] INFO  com.zaxxer.hikari.HikariDataSource -  {} - Hikari Handler DB Pool - Start completed.
2021-Jun-21 20:17:59 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Bound value [org.springframework.jdbc.datasource.ConnectionHolder@20a3e10c] for key [HikariDataSource (Hikari Handler DB Pool)] to thread [main]
2021-Jun-21 20:17:59 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Initializing transaction synchronization
2021-Jun-21 20:17:59 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Clearing transaction synchronization
2021-Jun-21 20:17:59 PM [main] DEBUG org.springframework.kafka.transaction.KafkaTransactionManager -  {} - Creating new transaction with name [com.example.SampleBaseApplication$Foo.sendToKafkaAndDB]: PROPAGATION_REQUIRED,ISOLATION_DEFAULT; 'chainedTxMang'
2021-Jun-21 20:17:59 PM [main] INFO  org.apache.kafka.clients.producer.ProducerConfig -  {} - ProducerConfig values: 
    acks = -1
    batch.size = 16384
    bootstrap.servers = [localhost:9093]
    buffer.memory = 33554432
    client.dns.lookup = default
    client.id = producer-transx-0
    compression.type = lz4
    connections.max.idle.ms = 540000
    delivery.timeout.ms = 120000
    enable.idempotence = true
    interceptor.classes = []
    key.serializer = class org.apache.kafka.common.serialization.StringSerializer
    linger.ms = 0
    max.block.ms = 60000
    max.in.flight.requests.per.connection = 5
    max.request.size = 1048576
    metadata.max.age.ms = 300000
    metadata.max.idle.ms = 300000
    metric.reporters = []
    metrics.num.samples = 2
    metrics.recording.level = INFO
    metrics.sample.window.ms = 30000
    partitioner.class = class org.apache.kafka.clients.producer.internals.DefaultPartitioner
    receive.buffer.bytes = 32768
    reconnect.backoff.max.ms = 1000
    reconnect.backoff.ms = 50
    request.timeout.ms = 30000
    retries = 1
    retry.backoff.ms = 100
    sasl.client.callback.handler.class = null
    sasl.jaas.config = null
    sasl.kerberos.kinit.cmd = /usr/bin/kinit
    sasl.kerberos.min.time.before.relogin = 60000
    sasl.kerberos.service.name = null
    sasl.kerberos.ticket.renew.jitter = 0.05
    sasl.kerberos.ticket.renew.window.factor = 0.8
    sasl.login.callback.handler.class = null
    sasl.login.class = null
    sasl.login.refresh.buffer.seconds = 300
    sasl.login.refresh.min.period.seconds = 60
    sasl.login.refresh.window.factor = 0.8
    sasl.login.refresh.window.jitter = 0.05
    sasl.mechanism = GSSAPI
    security.protocol = SSL
    security.providers = null
    send.buffer.bytes = 131072
    ssl.cipher.suites = null
    ssl.enabled.protocols = [TLSv1.2]
    ssl.endpoint.identification.algorithm = https
    ssl.key.password = [hidden]
    ssl.keymanager.algorithm = SunX509
    ssl.keystore.location = C:\Users\DC-USER\Downloads\sample\src\main\resources\kafka.server.keystore.jks
    ssl.keystore.password = [hidden]
    ssl.keystore.type = JKS
    ssl.protocol = TLSv1.2
    ssl.provider = null
    ssl.secure.random.implementation = null
    ssl.trustmanager.algorithm = PKIX
    ssl.truststore.location = C:\Users\DC-USER\Downloads\sample\src\main\resources\kafka.server.truststore.jks
    ssl.truststore.password = [hidden]
    ssl.truststore.type = JKS
    transaction.timeout.ms = 300000
    transactional.id = transx-0
    value.serializer = class org.apache.kafka.common.serialization.StringSerializer

2021-Jun-21 20:17:59 PM [main] INFO  org.apache.kafka.clients.producer.KafkaProducer -  {} - [Producer clientId=producer-transx-0, transactionalId=transx-0] Instantiated a transactional producer.
2021-Jun-21 20:18:00 PM [main] WARN  org.apache.kafka.clients.producer.ProducerConfig -  {} - The configuration 'The TransactionalId to use for transactional delivery. This enables reliability semantics which span multiple producer sessions since it allows the client to guarantee that transactions using the same TransactionalId have been completed prior to starting any new transactions. If no TransactionalId is provided, then the producer is limited to idempotent delivery. Note that <code>enable.idempotence</code> must be enabled if a TransactionalId is configured. The default is <code>null</code>, which means transactions cannot be used. Note that, by default, transactions require a cluster of at least three brokers which is the recommended setting for production; for development you can change this, by adjusting broker setting <code>transaction.state.log.replication.factor</code>.' was supplied but isn't a known config.
2021-Jun-21 20:18:00 PM [main] INFO  org.apache.kafka.common.utils.AppInfoParser -  {} - Kafka version: 2.5.0
2021-Jun-21 20:18:00 PM [main] INFO  org.apache.kafka.common.utils.AppInfoParser -  {} - Kafka commitId: 66563e712b0b9f84
2021-Jun-21 20:18:00 PM [main] INFO  org.apache.kafka.common.utils.AppInfoParser -  {} - Kafka startTimeMs: 1624286880090
2021-Jun-21 20:18:00 PM [main] INFO  org.apache.kafka.clients.producer.internals.TransactionManager -  {} - [Producer clientId=producer-transx-0, transactionalId=transx-0] Invoking InitProducerId for the first time in order to acquire a producer ID
2021-Jun-21 20:18:01 PM [kafka-producer-network-thread | producer-transx-0] INFO  org.apache.kafka.clients.Metadata -  {} - [Producer clientId=producer-transx-0, transactionalId=transx-0] Cluster ID: xfiOAAyKRtC9OMjttHLmyQ
2021-Jun-21 20:18:01 PM [kafka-producer-network-thread | producer-transx-0] INFO  org.apache.kafka.clients.producer.internals.TransactionManager -  {} - [Producer clientId=producer-transx-0, transactionalId=transx-0] Discovered transaction coordinator localhost:9093 (id: 1 rack: null)
2021-Jun-21 20:18:01 PM [kafka-producer-network-thread | producer-transx-0] INFO  org.apache.kafka.clients.producer.internals.TransactionManager -  {} - [Producer clientId=producer-transx-0, transactionalId=transx-0] ProducerId set to 0 with epoch 7
2021-Jun-21 20:18:01 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Bound value [org.springframework.kafka.core.KafkaResourceHolder@42aa1324] for key [org.springframework.kafka.core.DefaultKafkaProducerFactory@3976ebfa] to thread [main]
2021-Jun-21 20:18:01 PM [main] DEBUG org.springframework.kafka.transaction.KafkaTransactionManager -  {} - Created Kafka transaction on producer [CloseSafeProducer [delegate=org.apache.kafka.clients.producer.KafkaProducer@6164e137]]
2021-Jun-21 20:18:01 PM [main] TRACE org.springframework.transaction.interceptor.TransactionInterceptor -  {} - Getting transaction for [com.example.SampleBaseApplication$Foo.sendToKafkaAndDB]
2021-Jun-21 20:18:01 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Retrieved value [org.springframework.jdbc.datasource.ConnectionHolder@20a3e10c] for key [HikariDataSource (Hikari Handler DB Pool)] bound to thread [main]
2021-Jun-21 20:18:01 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Retrieved value [org.springframework.jdbc.datasource.ConnectionHolder@20a3e10c] for key [HikariDataSource (Hikari Handler DB Pool)] bound to thread [main]
2021-Jun-21 20:18:01 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Retrieved value [org.springframework.jdbc.datasource.ConnectionHolder@20a3e10c] for key [HikariDataSource (Hikari Handler DB Pool)] bound to thread [main]
2021-Jun-21 20:18:01 PM [main] INFO  com.example.DataAccess -  {} - Time taken to perform insert operation: 278ms
2021-Jun-21 20:18:01 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Retrieved value [org.springframework.kafka.core.KafkaResourceHolder@42aa1324] for key [org.springframework.kafka.core.DefaultKafkaProducerFactory@3976ebfa] bound to thread [main]
2021-Jun-21 20:18:01 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Retrieved value [org.springframework.kafka.core.KafkaResourceHolder@42aa1324] for key [org.springframework.kafka.core.DefaultKafkaProducerFactory@3976ebfa] bound to thread [main]
SendResult [producerRecord=ProducerRecord(topic=TEST_TOPIC, partition=null, headers=RecordHeaders(headers = [], isReadOnly = true), key=null, value=bar, timestamp=null), recordMetadata=TEST_TOPIC-0@0]
2021-Jun-21 20:18:02 PM [main] TRACE org.springframework.transaction.interceptor.TransactionInterceptor -  {} - Completing transaction for [com.example.SampleBaseApplication$Foo.sendToKafkaAndDB] after exception: java.lang.RuntimeException: exp...
2021-Jun-21 20:18:02 PM [main] TRACE org.springframework.transaction.interceptor.RuleBasedTransactionAttribute -  {} - Applying rules to determine whether transaction should rollback on java.lang.RuntimeException: exp...
2021-Jun-21 20:18:02 PM [main] TRACE org.springframework.transaction.interceptor.RuleBasedTransactionAttribute -  {} - Winning rollback rule is: null
2021-Jun-21 20:18:02 PM [main] TRACE org.springframework.transaction.interceptor.RuleBasedTransactionAttribute -  {} - No relevant rollback rule found: applying default rules
2021-Jun-21 20:18:02 PM [main] DEBUG org.springframework.kafka.transaction.KafkaTransactionManager -  {} - Initiating transaction rollback
2021-Jun-21 20:18:02 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Removed value [org.springframework.kafka.core.KafkaResourceHolder@42aa1324] for key [org.springframework.kafka.core.DefaultKafkaProducerFactory@3976ebfa] from thread [main]
2021-Jun-21 20:18:02 PM [main] DEBUG org.springframework.kafka.transaction.KafkaTransactionManager -  {} - Resuming suspended transaction after completion of inner transaction
2021-Jun-21 20:18:02 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Initializing transaction synchronization
2021-Jun-21 20:18:02 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Clearing transaction synchronization
2021-Jun-21 20:18:02 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Removed value [org.springframework.jdbc.datasource.ConnectionHolder@20a3e10c] for key [HikariDataSource (Hikari Handler DB Pool)] from thread [main]
2021-Jun-21 20:18:02 PM [main] TRACE org.springframework.transaction.support.TransactionSynchronizationManager -  {} - Initializing transaction synchronization
2021-Jun-21 20:18:02 PM [main] INFO  org.springframework.boot.autoconfigure.logging.ConditionEvaluationReportLoggingListener -  {} - 

Error starting ApplicationContext. To display the conditions report re-run your application with 'debug' enabled.
2021-Jun-21 20:18:02 PM [main] ERROR org.springframework.boot.SpringApplication -  {} - Application run failed
java.lang.IllegalStateException: Failed to execute ApplicationRunner
    at org.springframework.boot.SpringApplication.callRunner(SpringApplication.java:789) [spring-boot-2.3.1.RELEASE.jar:2.3.1.RELEASE]
    at org.springframework.boot.SpringApplication.callRunners(SpringApplication.java:776) [spring-boot-2.3.1.RELEASE.jar:2.3.1.RELEASE]
    at org.springframework.boot.SpringApplication.run(SpringApplication.java:322) [spring-boot-2.3.1.RELEASE.jar:2.3.1.RELEASE]
    at org.springframework.boot.SpringApplication.run(SpringApplication.java:1237) [spring-boot-2.3.1.RELEASE.jar:2.3.1.RELEASE]
    at org.springframework.boot.SpringApplication.run(SpringApplication.java:1226) [spring-boot-2.3.1.RELEASE.jar:2.3.1.RELEASE]
    at com.example.SampleBaseApplication.main(SampleBaseApplication.java:60) [classes/:?]
Caused by: java.lang.RuntimeException: exp...
    at com.example.SampleBaseApplication$Foo.sendToKafkaAndDB(SampleBaseApplication.java:102) ~[classes/:?]
    at com.example.SampleBaseApplication$Foo$$FastClassBySpringCGLIB$$ca642166.invoke(<generated>) ~[classes/:?]
    at org.springframework.cglib.proxy.MethodProxy.invoke(MethodProxy.java:218) ~[spring-core-5.2.7.RELEASE.jar:5.2.7.RELEASE]
    at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.invokeJoinpoint(CglibAopProxy.java:771) ~[spring-aop-5.2.7.RELEASE.jar:5.2.7.RELEASE]
    at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:163) ~[spring-aop-5.2.7.RELEASE.jar:5.2.7.RELEASE]
    at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.proceed(CglibAopProxy.java:749) ~[spring-aop-5.2.7.RELEASE.jar:5.2.7.RELEASE]
    at org.springframework.transaction.interceptor.TransactionInterceptor$$Lambda$805/1148088421.proceedWithInvocation(Unknown Source) ~[?:?]
    at org.springframework.transaction.interceptor.TransactionAspectSupport.invokeWithinTransaction(TransactionAspectSupport.java:367) ~[spring-tx-5.2.7.RELEASE.jar:5.2.7.RELEASE]
    at org.springframework.transaction.interceptor.TransactionInterceptor.invoke(TransactionInterceptor.java:118) ~[spring-tx-5.2.7.RELEASE.jar:5.2.7.RELEASE]
    at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:186) ~[spring-aop-5.2.7.RELEASE.jar:5.2.7.RELEASE]
    at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.proceed(CglibAopProxy.java:749) ~[spring-aop-5.2.7.RELEASE.jar:5.2.7.RELEASE]
    at org.springframework.aop.framework.CglibAopProxy$DynamicAdvisedInterceptor.intercept(CglibAopProxy.java:691) ~[spring-aop-5.2.7.RELEASE.jar:5.2.7.RELEASE]
    at com.example.SampleBaseApplication$Foo$$EnhancerBySpringCGLIB$$9728c463.sendToKafkaAndDB(<generated>) ~[classes/:?]
    at com.example.SampleBaseApplication.lambda$0(SampleBaseApplication.java:68) ~[classes/:?]
    at com.example.SampleBaseApplication$$Lambda$526/1532644077.run(Unknown Source) ~[?:?]
    at org.springframework.boot.SpringApplication.callRunner(SpringApplication.java:786) ~[spring-boot-2.3.1.RELEASE.jar:2.3.1.RELEASE]
    ... 5 more
2021-Jun-21 20:18:02 PM [main] INFO  org.springframework.boot.web.embedded.tomcat.GracefulShutdown -  {} - Commencing graceful shutdown. Waiting for active requests to complete
2021-Jun-21 20:18:02 PM [tomcat-shutdown] INFO  org.springframework.boot.web.embedded.tomcat.GracefulShutdown -  {} - Graceful shutdown complete
2021-Jun-21 20:18:02 PM [main] INFO  org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor -  {} - Shutting down ExecutorService 'applicationTaskExecutor'
2021-Jun-21 20:18:02 PM [main] INFO  org.apache.kafka.clients.producer.KafkaProducer -  {} - [Producer clientId=producer-transx-0, transactionalId=transx-0] Closing the Kafka producer with timeoutMillis = 30000 ms.
2021-Jun-21 20:18:02 PM [main] INFO  com.zaxxer.hikari.HikariDataSource -  {} - Hikari Handler DB Pool - Shutdown initiated...
2021-Jun-21 20:18:02 PM [main] INFO  com.zaxxer.hikari.HikariDataSource -  {} - Hikari Handler DB Pool - Shutdown completed.
0 Answers
Related