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.