【问题标题】:Implementation of a step-scoped in-memory destination in spring batch在 Spring Batch 中实现步进范围的内存中目标
【发布时间】:2015-08-04 17:46:44
【问题描述】:

在使用 Spring Batch 创建大型报告时,我有一个要求,以二进制格式保存在数据库中。任何工作数据都不能直接写入文件或JobExecutionContext之外的工作表。

我知道通常你只会写作业执行上下文,但我有点困惑如何处理这么大的报告(可能有几百兆字节)。

目前,我的 Writer 实现依赖于一个聚合器类,该类作为 bean 注入,然后有一个 TaskLet 注入了聚合器,将完成的报告写入数据库。

问题是我无法将聚合器的范围限定为 step 上下文,因此如果两个作业同时运行,它们将写入同一个聚合器。

这是我当前的实现

域类

public class DataChunk {
    private int pageNumber;
    private byte[] data;
}

作家

public class FooWriter implements ItemWriter<DataChunk> {

    private DataChunkAggregator dataChunkAggregator;

    public void write(List<? extends DataChunk> dataChunks) throws Exception {
        dataChunks.stream().forEach(chunk -> dataChunkAggregator.addChunk(chunk.getPageNumber(), chunk.getData()));
    }
}

聚合器

public class FooAggregator {
    private Map<int, byte> pagedData; // Key sorted implementation

    public void addChunk(int pageNumber, byte[] data) {
        pagedData.put(pageNumber, data)
    }

    public byte[] aggregate() {
        ByteArrayOutputStream baos = new ByteArrayOutputStream();
        pagedData.values.stream().forEach(data -> baos.write(data));
        return baos.toByteArray();
    }
}

报告写作小任务

public class ReportWritingTasklet implements TaskLet {

    private ReportRepository reportRepository;
    private FooAggregator fooAggregator;

    public RepeatStatus execute(StepContribution contribution, ChunkContext context) {
        byte[] data = fooAggregator.aggregate();
        reportRepository.getOne(reportId).setDataBytes(data);
    }
}

上下文

<?xml version="1.0" encoding="UTF-8"?>
<beans>
    <bean id=fooWriter class="FooWriter" scope="step"
        p:fooAggregator-ref="fooAggregator"/>

    <bean id="fooAggregator" class="FooAggregator"/>

    <bean id="reportWritingTasklet" class="ReportWritingTasklet" scope="step"
        p:fooAggregator-ref="fooAggregator"/>

    <batch:job id="fooJob">
        <batch:step id="generateReport" next="assembleReport">
            <batch:chunk reader="fooReader" processor="fooProcessor" writer="fooWriter"/>
        </batch:step>
        <batch:step id="assembleReport">
            <batch:tasklet class="ReportWritingTasklet"/>
        </batch:step>
     </batch:job>
</beans>

如果我尝试将 FooAggregator 设为步进范围,我会得到以下异常作为根本原因

Caused by: java.lang.IllegalStateException: Cannot convert value of type    [com.sun.proxy.$Proxy98 implementing org.springframework.aop.scope.ScopedObject,java.io.Serializable,org.springframework.aop.framework.AopInfrastructureBean,org.springframework.aop.SpringProxy,org.springframework.aop.framework.Advised] to required type [FooAggregator] for property 'fooAggregator': no matching editors or conversion strategy found

这是因为您只能将某些内容限定为该步骤。

如何将执行上下文用作我的数据块的接收器,记住它们会很多而且它们会非常大?

【问题讨论】:

    标签: java spring java-8 spring-batch batch-processing


    【解决方案1】:

    我已经设法解决了这个问题。它不是很 Spring batch-y,但它符合我的要求。

    本质上,有太多数据需要进出上下文。解决方案是保持编写器本身的状态,并通过TransactionCallback 将其保存在Step 末尾的StepExecutionListener

    更新的Writer

    public class FooWriter extends StepExecutionListenerSupport implements ItemWriter<DataChunk> {
    
        private String reportId;
        private Map<Integer, byte[]> byteArrayMap = new ConcurrentSkipListMap<>();
    
        private TransactionTemplate transactionTemplate;
        private ReportRepository reportRepository;
    
        @Override
        public synchronized void write(List<? extends CaseChunk> caseChunks) throws Exception {
            caseChunks.stream().forEach(chunk -> {
                byteArrayMap.put(chunk.getPageNumber(), chunk.getBytes());
            });
        }
    
        @Override
        public void beforeStep(StepExecution stepExecution) {
            // No-op
        }
    
        @Override
        public ExitStatus afterStep(StepExecution stepExecution) {
            StringBuilder sb = new StringBuilder();
            for (byte[] byteArrayOutputStream : byteArrayMap.values()) {
                sb.append(new String(Base64.decode(byteArrayOutputStream)));
            }
            String encodedReportData = new String(Base64.encode(sb.toString().getBytes()));
    
            TransactionCallback<Report> transactionCallback = transactionStatus -> {
                Report report = reportRepository.getOne(this.reportId);
                report.setReportData(encodedReportData);
                reportRepository.save(report);
                return report;
            };
    
            // TransactionTemplate throws its own declared TransactionException, rethrows encountered RuntimeExceptions
            // and also Errors. Any problem writing the date kills the job, so it's OK to catch Throwable here instead
            // of trying to
            try {
                transactionTemplate.execute(transactionCallback);
            } catch (Throwable t) {
                LOGGER.error("Error saving report data ID:[{}]", reportId);
                return ExitStatus.FAILED.addExitDescription(t);
            }
            return ExitStatus.COMPLETED;
        }
    
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-01-04
      • 2019-01-04
      • 2013-02-10
      • 2016-07-13
      • 2011-01-31
      • 1970-01-01
      • 2012-04-11
      • 1970-01-01
      相关资源
      最近更新 更多