【问题标题】:Can Same spring batch invoked with different job parameter from different multiple thread at the same time可以同时从不同的多个线程使用不同的作业参数调用相同的弹簧批处理
【发布时间】: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


    【解决方案1】:

    MapJobRepository用于生产用途。它是线程安全的。如果您需要内存作业存储库的性能(失去可重新启动性等),请使用内存数据库,如 HSQLDB。

    除此之外,如果您使用线程安全组件,则没有理由不能启动具有多个线程的多个作业实例。

    【讨论】:

      【解决方案2】:

      您确定 MapJobExecutionDao 在所有方面都是线程安全的吗?我看到,在 MapJobExecutionDao 中使用了 ConcurrentMap,但我不确定这是否足够。我曾经在从不同线程访问的 Map 中获取 NullPointer 时遇到问题。问题是,一个线程引起了重新散列,当第二个线程在那一刻确实访问了地图时,它收到了一个空指针。

      您确定您的识别作业参数的组合是独一无二的吗?我明白了,您使用 System.currentTimeMillis() 添加了一个参数 Job_Time,但是您知道这是否真的以唯一的时间戳解决?

      您是否尝试过使用基于表格的 JobExecutionDao 等版本?

      【讨论】:

        猜你喜欢
        • 2021-11-15
        • 1970-01-01
        • 2014-07-29
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-07-20
        • 1970-01-01
        相关资源
        最近更新 更多