From aef9a52d402ce47956405e669945f2c394f7f14a Mon Sep 17 00:00:00 2001 From: Michael Minella Date: Mon, 9 Nov 2015 13:36:57 -0600 Subject: [PATCH] Added test to verify concurrent transactions with HSQLDB --- build.gradle | 10 +- .../ConcurrentTransactionTests.java | 219 ++++++++++++++++++ 2 files changed, 225 insertions(+), 4 deletions(-) create mode 100644 spring-batch-core-tests/src/test/java/org/springframework/batch/core/test/concurrent/ConcurrentTransactionTests.java diff --git a/build.gradle b/build.gradle index fbdc5d0ec..93b951cc2 100644 --- a/build.gradle +++ b/build.gradle @@ -359,6 +359,7 @@ project('spring-batch-core-tests') { project.tasks.findByPath("artifactoryPublish")?.enabled = false dependencies { compile project(":spring-batch-core") + compile project(":spring-batch-infrastructure") compile "commons-dbcp:commons-dbcp:$commonsDdbcpVersion" compile "org.springframework:spring-jdbc:$springVersion" compile "org.springframework.retry:spring-retry:$springRetryVersion" @@ -368,12 +369,13 @@ project('spring-batch-core-tests') { testCompile "org.hsqldb:hsqldb:$hsqldbVersion" testCompile "commons-io:commons-io:$commonsIoVersion" testCompile "org.apache.derby:derby:$derbyVersion" - testCompile("junit:junit:${junitVersion}") { - exclude group:'org.hamcrest', module:'hamcrest-core' - } - testCompile("org.hamcrest:hamcrest-all:$hamcrestVersion") + testCompile("junit:junit:${junitVersion}") { + exclude group:'org.hamcrest', module:'hamcrest-core' + } + testCompile("org.hamcrest:hamcrest-all:$hamcrestVersion") testCompile "log4j:log4j:$log4jVersion" testCompile "org.springframework:spring-test:$springVersion" + testCompile "org.springframework:spring-jdbc:$springVersion" runtime "mysql:mysql-connector-java:$mysqlVersion" runtime "postgresql:postgresql:$postgresqlVersion" diff --git a/spring-batch-core-tests/src/test/java/org/springframework/batch/core/test/concurrent/ConcurrentTransactionTests.java b/spring-batch-core-tests/src/test/java/org/springframework/batch/core/test/concurrent/ConcurrentTransactionTests.java new file mode 100644 index 000000000..bf21824b8 --- /dev/null +++ b/spring-batch-core-tests/src/test/java/org/springframework/batch/core/test/concurrent/ConcurrentTransactionTests.java @@ -0,0 +1,219 @@ +/* + * Copyright 2015 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.batch.core.test.concurrent; + +import static org.junit.Assert.assertEquals; + +import java.sql.Connection; +import java.sql.Driver; +import java.sql.SQLException; +import java.sql.Statement; + +import javax.sql.DataSource; + +import org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.batch.core.BatchStatus; +import org.springframework.batch.core.Job; +import org.springframework.batch.core.JobExecution; +import org.springframework.batch.core.JobParameters; +import org.springframework.batch.core.Step; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.configuration.annotation.DefaultBatchConfigurer; +import org.springframework.batch.core.configuration.annotation.EnableBatchProcessing; +import org.springframework.batch.core.configuration.annotation.JobBuilderFactory; +import org.springframework.batch.core.configuration.annotation.StepBuilderFactory; +import org.springframework.batch.core.job.builder.FlowBuilder; +import org.springframework.batch.core.job.flow.Flow; +import org.springframework.batch.core.launch.JobLauncher; +import org.springframework.batch.core.repository.JobRepository; +import org.springframework.batch.core.repository.support.JobRepositoryFactoryBean; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.core.io.DefaultResourceLoader; +import org.springframework.core.io.ResourceLoader; +import org.springframework.core.task.SimpleAsyncTaskExecutor; +import org.springframework.core.task.TaskExecutor; +import org.springframework.jdbc.datasource.embedded.ConnectionProperties; +import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseConfigurer; +import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseFactory; +import org.springframework.jdbc.datasource.init.ResourceDatabasePopulator; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.util.ClassUtils; + +/** + * @author Michael Minella + */ +@RunWith(SpringJUnit4ClassRunner.class) +@ContextConfiguration(classes = ConcurrentTransactionTests.ConcurrentJobConfiguration.class) +public class ConcurrentTransactionTests { + + @Autowired + private Job concurrentJob; + + @Autowired + private JobLauncher jobLauncher; + + @DirtiesContext + @Test + public void testConcurrentLongRunningJobExecutions() throws Exception { + + JobExecution jobExecution = jobLauncher.run(concurrentJob, new JobParameters()); + + assertEquals(jobExecution.getStatus(), BatchStatus.COMPLETED); + } + + @Configuration + @EnableBatchProcessing + public static class ConcurrentJobConfiguration extends DefaultBatchConfigurer { + + @Autowired + private JobBuilderFactory jobBuilderFactory; + + @Autowired + private StepBuilderFactory stepBuilderFactory; + + @Bean + public TaskExecutor taskExecutor() { + return new SimpleAsyncTaskExecutor(); + } + + /** + * This datasource configuration configures the HSQLDB instance using MVCC. When + * configurd using the default behavior, transaction serialization errors are + * thrown (default configuration example below). + * + * return new PooledEmbeddedDataSource(new EmbeddedDatabaseBuilder(). + * addScript("classpath:org/springframework/batch/core/schema-drop-hsqldb.sql"). + * addScript("classpath:org/springframework/batch/core/schema-hsqldb.sql"). + * build()); + + * @return + */ + @Bean + DataSource dataSource() { + ResourceLoader defaultResourceLoader = new DefaultResourceLoader(); + EmbeddedDatabaseFactory embeddedDatabaseFactory = new EmbeddedDatabaseFactory(); + embeddedDatabaseFactory.setDatabaseConfigurer(new EmbeddedDatabaseConfigurer() { + + @Override + public void configureConnectionProperties(ConnectionProperties properties, String databaseName) { + try { + properties.setDriverClass((Class) ClassUtils.forName("org.hsqldb.jdbcDriver", this.getClass().getClassLoader())); + } + catch (Exception e) { + e.printStackTrace(); + } + properties.setUrl("jdbc:hsqldb:mem:" + databaseName + ";hsqldb.tx=mvcc"); + properties.setUsername("sa"); + properties.setPassword(""); + } + + @Override + public void shutdown(DataSource dataSource, String databaseName) { + try { + Connection connection = dataSource.getConnection(); + Statement stmt = connection.createStatement(); + stmt.execute("SHUTDOWN"); + } + catch (SQLException ex) { + } + } + }); + + ResourceDatabasePopulator databasePopulator = new ResourceDatabasePopulator(); + databasePopulator.addScript(defaultResourceLoader.getResource("classpath:org/springframework/batch/core/schema-drop-hsqldb.sql")); + databasePopulator.addScript(defaultResourceLoader.getResource("classpath:org/springframework/batch/core/schema-hsqldb.sql")); + embeddedDatabaseFactory.setDatabasePopulator(databasePopulator); + + return embeddedDatabaseFactory.getDatabase(); + } + + @Bean + public Flow flow() { + return new FlowBuilder("flow") + .start(stepBuilderFactory.get("flow.step1") + .tasklet(new Tasklet() { + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { + return RepeatStatus.FINISHED; + } + }).build() + ).next(stepBuilderFactory.get("flow.step2") + .tasklet(new Tasklet() { + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { + return RepeatStatus.FINISHED; + } + }).build() + ).build(); + } + + @Bean + public Step firstStep() { + return stepBuilderFactory.get("firstStep") + .tasklet(new Tasklet() { + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { + System.out.println(">> Beginning concurrent job test"); + return RepeatStatus.FINISHED; + } + }).build(); + } + + @Bean + public Step lastStep() { + return stepBuilderFactory.get("lastStep") + .tasklet(new Tasklet() { + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { + System.out.println(">> Ending concurrent job test"); + return RepeatStatus.FINISHED; + } + }).build(); + } + + @Bean + public Job concurrentJob() { + Flow splitFlow = new FlowBuilder("splitflow").split(new SimpleAsyncTaskExecutor()).add(flow(), flow(), flow(), flow(), flow(), flow(), flow()).build(); + + return jobBuilderFactory.get("concurrentJob") + .start(firstStep()) + .next(stepBuilderFactory.get("splitFlowStep") + .flow(splitFlow) + .build()) + .next(lastStep()) + .build(); + } + + @Override + protected JobRepository createJobRepository() throws Exception { + JobRepositoryFactoryBean factory = new JobRepositoryFactoryBean(); + factory.setDataSource(dataSource()); + factory.setIsolationLevelForCreate("ISOLATION_READ_COMMITTED"); + factory.setTransactionManager(getTransactionManager()); + factory.afterPropertiesSet(); + return factory.getObject(); + } + } +}