【问题标题】:How to make multiple writer execute parallely in spring batch如何使多个编写器在春季批处理中并行执行
【发布时间】:2021-12-28 19:34:19
【问题描述】:

下面的代码是我现有的代码,它一个接一个地处理和写入数据到多个集合中。

我的要求是我想同时写入多个集合,而不是一个接一个。总之我想做并行写入过程

下面是我现有的代码

    public Job importSingleETLData(SingleETLJobListener listener, Step singleETLStep, HttpServletRequest request) {
        return jobBuilderFactory.get("importSingleETLData")
                .incrementer(new RunIdIncrementer())
                .listener(listener)
                .flow(singleETLStep)
                .end()
                .build();
    }


@Bean
    public Step singleETLStep(MongoItemWriter<CompositeWriterData> writer, HttpServletRequest request) {
        return stepBuilderFactory.get("singleETLStep")
                // TODO: P3 chunk size configurable
                .<UserInfo, CompositeWriterData>chunk(etlConfiguration.getBatchChunkSize())
                .reader(reader(("#{jobParameters[profileId]}"))).faultTolerant().skipPolicy(readerSkipper())
                .processor(processor(request,"#{jobParameters[executeProcessing]}"))
                .listener(processorListener()).faultTolerant().skipPolicy(writerSkipper())
                .writer(writer)
                .build();
    }```


    ```@Bean
    @StepScope
    public MongoItemReader<UserInfo> reader(@Value("#{jobParameters[profileId]}") String profileId) {

        String query = "{'results._id' :'"+profileId + "'}";
        Map<String, Direction> sorts = new HashMap<>();
        sorts.put("_id", Direction.ASC);
        MongoItemReader<UserInfo> reader = new MongoItemReader<>();
        reader.setCollection(CommonConstants.USER_INFO_VIEW);
        reader.setTemplate(secondaryMongoTemplate);
        reader.setTargetType(UserInfo.class);
        // TODO: P2 take latest phi only
        // TODO: P2 Use different query in on demand to fetch last processed record if nothing is updated from when last process ran 
        reader.setQuery(query);
        reader.setSort(sorts);
        return reader;
    }

    @Bean
    @StepScope
    public ETLDataProcessor processor(HttpServletRequest request,@Value("#{jobParameters[executeProcessing]}") String executeProcessing) {
        return new ETLDataProcessor(request,executeProcessing);
    }

    @Bean
    public MongoItemWriter<ProfileRecommendationInfo> recommendationsDataWriter() {
        MongoItemWriter<ProfileRecommendationInfo> writer = new MongoItemWriter<>();
        writer.setTemplate(secondaryMongoTemplate);
        return writer;
    }

    @Bean
    public MongoItemWriter<ProfileLifebandInfo> lifeBandDataWriter() {
        MongoItemWriter<ProfileLifebandInfo> writer = new MongoItemWriter<>();
        writer.setTemplate(secondaryMongoTemplate);
        return writer;
    }

    @Bean
    public MongoItemWriter<Profile> profileWriter() {
        MongoItemWriter<Profile> writer = new MongoItemWriter<>();
        writer.setTemplate(secondaryMongoTemplate);
        return writer;
    }

    
    @Bean
    public MongoItemWriter<CompositeWriterData> compositeMongoWriter() {
        CompositeMongoItemWriter compositeWriter = new CompositeMongoItemWriter();
        compositeWriter.setTemplate(secondaryMongoTemplate);
        return compositeWriter;
    }

    @Bean
    public SingleETLProcessorListener processorListener() {
        return new SingleETLProcessorListener();
    }

    @Bean
    public SkipPolicy readerSkipper() {
        return new ReaderSkipper();
    }
    
    @Bean
    public SkipPolicy writerSkipper() {
        return new WriterSkipper();
    }
public class CompositeMongoItemWriter extends MongoItemWriter<CompositeWriterData> {

    @Autowired
    MongoItemWriter<ProfileRecommendationInfo> recommendationsDataWriter;
    @Autowired
    MongoItemWriter<ProfileLifebandInfo> lifeBandWriter;
    @Autowired
    private MongoTemplate secondaryMongoTemplate;
    @Autowired
    MongoItemWriter<Profile> profileWriter;

    @Override
    public void write(List<? extends CompositeWriterData> items) throws Exception {
        if( items!= null && !items.isEmpty()) {

            for(CompositeWriterData compositeWriterData : items) {
                for( Entry<String, Object> collection : compositeWriterData.getCollectionsPOJODataMap().entrySet() ) {

                    MongoItemWriter mongoItemWriter = fetchMongoItemWriterObject(collection.getKey());

                    if(mongoItemWriter != null) {
                        mongoItemWriter.write(Arrays.asList(collection.getValue()));
                    }

                    // Below code will update Profile with profile_recommendation_id.primary key
                    if(CommonConstants.PROFILE_RECOMMENDATION_INFO.equals(collection.getKey())) {
                        ProfileRecommendationInfo profileRecommendationInfo = (ProfileRecommendationInfo) collection.getValue();
                        updateProfileWithSavedCollectionDataId(profileRecommendationInfo.getProfileId(),profileRecommendationInfo.getDataId());
                    }

                }

            }
        }
    }

    /**
     *  This method will return an object of MongoItemWriter based on the collectionName
     * passed to it on invocation.
     * 
     * Note: - This method needs to be modified whenever new collection is added in
     * ${etl.processor.collection.pojo} in utility-service-application.properties
     * @return
     */
    private MongoItemWriter fetchMongoItemWriterObject(String collectionName){
        if(CommonConstants.PROFILE_RECOMMENDATION_INFO.equalsIgnoreCase(collectionName)) {
            return recommendationsDataWriter;
        }else if(CommonConstants.PROFILE_LIFEBAND_INFO.equalsIgnoreCase(collectionName)) {
            return lifeBandWriter;
        }
        return null;
    }

    /**
     *  Below method will update profile. with the collection data primary key value
     *          This is useful in etl processing when we fetch last data saved for a particular user
     * @throws Exception 
     */
    private void updateProfileWithSavedCollectionDataId(String profileId,String profileRecommendationInfoId) throws Exception {
        Profile profile = secondaryMongoTemplate.findById(profileId, Profile.class);
        profile.setProfileRecommendationInfoId(profileRecommendationInfoId);
        profileWriter.write(Arrays.asList(profile));
    }
}

如何同时写入多个集合而不是一个接一个。总之我想做并行写入过程

我们正在努力实现以下链接中给出的并行处理

https://docs.spring.io/spring-batch/docs/current/reference/html/scalability.html

【问题讨论】:

    标签: spring-boot spring-batch batch-processing


    【解决方案1】:

    CompositeItemWriter 按顺序调用委托编写者。如果您想并行调用委托编写器,您需要一个自定义的CompositeIemWriter,例如将不同的写入操作提交给TaskExecutor(参见示例here)。但是,您需要考虑错误处理以及如何从中恢复(即重试/跳过功能)。

    【讨论】:

    • 谢谢@Mahmoud Ben Hassine,会检查
    猜你喜欢
    • 2018-02-08
    • 1970-01-01
    • 2020-03-31
    • 2014-06-30
    • 2019-03-13
    • 1970-01-01
    • 2016-03-03
    • 2019-03-26
    • 1970-01-01
    相关资源
    最近更新 更多