【发布时间】:2015-10-08 04:37:40
【问题描述】:
我们的要求是同时写入多个文件。我们正在使用弹簧批处理来编写文件,并且我们正在从不同的线程中处理弹簧批处理。每个线程都有自己的应用程序上下文。所以我们可以保证单例bean不会在多个线程之间共享。下面是我的代码 sn-p。
Spring 批处理配置。
<bean id="reportDataReader" class="com.test.ist.batch2.rrm.batch.readers.RRMItmeReader"
scope="step">
<property name="verifyCursorPosition" value="false" />
<property name="dataSource" ref="dataSource" />
<property name="sql" value="#{jobParameters['sqlquery']}" />
<property name="rowMapper" ref="valueMapper" />
<property name="fetchSize" value="5000" />
</bean>
<bean id="valueMapper" class="com.test.ist.batch2.rrm.batch.mappers.DBValueMapper" scope="step"></bean>
<bean id="velocityFileWritter"
class="com.test.ist.batch2.rrm.batch.writers.RRMVelocityFileWriter"
scope="step">
</bean>
<bean id="velocityEngine"
class="org.springframework.ui.velocity.VelocityEngineFactoryBean">
<property name="velocityProperties">
<value>
resource.loader = class
class.resource.loader.class = org.apache.velocity.runtime.resource.loader.ClasspathResourceLoader
class.resource.loader.cache = true
class.resource.loader.modificationCheckInterval = 0
</value>
</property>
</bean>
<batch:job id="rrmReportGenJob">
<batch:step id="rrmReportGenStep">
<batch:tasklet>
<batch:chunk reader="reportDataReader" writer="velocityFileWritter"
commit-interval="${reportData.reader.commit-interval}">
</batch:chunk>
</batch:tasklet>
</batch:step>
</batch:job>
这就是我们调用 spring 批处理的方式。
ThreadPoolExecutor tpe=new ThreadPoolExecutor(10, 10, 1000000, TimeUnit.MILLISECONDS, new LinkedBlockingQueue()); PetReportGenerator rrg=new PetReportGenerator(null); ThreadTest tt=new ThreadTest(new PetReportGenerator(null), "161"); ThreadTest tt2=new ThreadTest(new PetReportGenerator(null), "162"); ThreadTest tt3=new ThreadTest(new PetReportGenerator(null), "163"); ThreadTest tt4=new ThreadTest(new PetReportGenerator(null), "165"); tpe.execute(tt); tpe.execute(tt2); tpe.execute(tt3); tpe.execute(tt4);
在 PetReportGenerator 的构造函数中,我们正在初始化 bean 配置。 下面是代码sn-p
私有 ApplicationContext appContext;
public PetReportGenerator(ApplicationContext reportContext){
if(null == reportContext){
//if(null == appContext){
appContext=new ClassPathXmlApplicationContext("spring-batch-jobs.xml");
//}
}else{
setAppContext(reportContext);;
}
}
下面是我们如何调用spring批处理的代码摘录
Job jobToExecute = (Job)SpringUtils.getBean(jobName); JobParametersBuilder paramsBuilder = new JobParametersBuilder(); //默认添加数据时间。这将有助于使用相同的参数再次启动相同的作业 paramsBuilder.addLong("JOB_TIME", System.currentTimeMillis()); 如果(!jobParams.isEmpty()){ //验证输入字段。 String sqlToUse = validator.validateInput(jobParams); for(Map.Entry entry:jobParams.entrySet()){
paramsBuilder.addString(entry.getKey(), entry.getValue());
}
}else{
throw new ReportGenerationException("Job input parameter is Empty");
}
jobexe=jobLauncher.run(jobToExecute, paramsBuilder.toJobParameters());
如果它在单个线程中运行,它工作正常。 当它被多个线程调用时,我们会遇到错误
09:09:26,742 ERROR pool-1-thread-3 job.AbstractJob:329 - 执行作业时遇到致命错误 java.lang.NullPointerException 在 org.springframework.batch.core.repository.dao.MapJobExecutionDao.synchronizeStatus(MapJobExecutionDao.java:158) 在 org.springframework.batch.core.repository.support.SimpleJobRepository.update(SimpleJobRepository.java:161) 在 sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 在 sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57) 在 sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) 在 java.lang.reflect.Method.invoke(Method.java:606) 在 org.springframework.aop.support.AopUtils.invokeJoinpointUsingReflection(AopUtils.java:317) 在 org.springframework.aop.framework.ReflectiveMethodInvocation.invokeJoinpoint(ReflectiveMethodInvocation.java:190) 在 org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:157) 在 org.springframework.transaction.interceptor.TransactionInterceptor$1.proceedWithInvocation(TransactionInterceptor.java:98) 在 org.springframework.transaction.interceptor.TransactionAspectSupport.invokeWithinTransaction(TransactionAspectSupport.java:262) 在 org.springframework.transaction.interceptor.TransactionInterceptor.invoke(TransactionInterceptor.java:95) 在 org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:179) 在 org.springframework.aop.framework.JdkDynamicAopProxy.invoke(JdkDynamicAopProxy.java:207) 在 com.sun.proxy.$Proxy14.update(未知来源) 在 org.springframework.batch.core.job.AbstractJob.updateStatus(AbstractJob.java:416) 在 org.springframework.batch.core.job.AbstractJob.execute(AbstractJob.java:299) 在 org.springframework.batch.core.launch.support.SimpleJobLauncher$1.run(SimpleJobLauncher.java:135) 在 org.springframework.core.task.SyncTaskExecutor.execute(SyncTaskExecutor.java:50) 在 org.springframework.batch.core.launch.support.SimpleJobLauncher.run(SimpleJobLauncher.java:128)
谁能帮我理解可能是什么问题?
【问题讨论】:
标签: spring-batch