【问题标题】:How to split a file using FileSplitter (@Splitter) in Spring Integration using java configuration如何使用 java 配置在 Spring Integration 中使用 FileSplitter (@Splitter) 拆分文件
【发布时间】:2018-08-06 08:44:48
【问题描述】:

我想根据 '\n' 拆分我的消息。拆分后,我想将拆分消息填充为 1000 个块进行处理。有没有办法在不使用迭代器或循环的情况下做到这一点? 我们也可以使用一些东西作为序列大小吗?以下是使用的格式

@InboundChannelAdaptor(channel = "x", poller = @Poller(fixedDelay = "20000", maxMessagesPerPoll = "1")

// 这个Channel用来获取文件

@Splitter(inputChannel = "x")

//这里我想使用 FileSplitter 将消息拆分成块 //fileSplitter.setOutputChannelName("y")

@ServiceActivator(inputChannel = "y")

//分块处理逻辑

更新 ----------实现-----------

@Bean
@InboundChannelAdapter(channel = "fileInputChannel", poller = @Poller(fixedDelay = "5000"))
public MessageSource<File> sftpMessageSource() {
    FileReadingMessageSource source = new FileReadingMessageSource();
    source.setDirectory(new File(INBOUND_PATH));
    source.setFilter(new AcceptOnceFileListFilter<>());
    return source;
}

@Splitter(inputChannel = "fileInputChannel")
@Bean
public FileSplitter fileSplitter() {
   FileSplitter fileSplitter = new FileSplitter();
   fileSplitter.setOutputChannelName("chunkingChannel");
   return fileSplitter;
}

@ServiceActivator(inputChannel = "chunkingChannel")
@Bean
public AggregatingMessageHandler chunker() {
    AggregatingMessageHandler aggregator = new AggregatingMessageHandler(new DefaultAggregatingMessageGroupProcessor());
    aggregator.setReleaseStrategy(new MessageCountReleaseStrategy(1000));
    aggregator.setExpireGroupsUponCompletion(true);
    aggregator.setGroupTimeoutExpression(new ValueExpression<>(100L));
    aggregator.setSendPartialResultOnExpiry(true);
    aggregator.setOutputChannelName("processFileChannel");
    return aggregator;
}

@Bean
@ServiceActivator(inputChannel = "processFileChannel")
public MessageHandler handler() {
    return new MessageHandler() {

        @Override
        public void handleMessage(Message<?> message) throws MessagingException {
            List<String> strings = (List<String>) message.getPayload();
            System.out.println( "List Size :  "+ strings.size() + " for List " + strings.toString());
        }

    };
}

我在使用 AggregatorFactoryBean 时遇到的错误

   org.springframework.beans.factory.BeanCreationException: Error creating bean with name 'datastreamApplication': Initialization of bean failed; nested exception is org.springframework.beans.factory.BeanCreationException: Error creating bean with name 'chunker': FactoryBean threw exception on object creation; nested exception is java.lang.IllegalArgumentException: targetObject must not be null
    at org.springframework.beans.factory.support.AbstractAutowireCapableBeanFactory.doCreateBean(AbstractAutowireCapableBeanFactory.java:581) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.AbstractAutowireCapableBeanFactory.createBean(AbstractAutowireCapableBeanFactory.java:495) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.AbstractBeanFactory.lambda$doGetBean$0(AbstractBeanFactory.java:317) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.DefaultSingletonBeanRegistry.getSingleton(DefaultSingletonBeanRegistry.java:222) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.AbstractBeanFactory.doGetBean(AbstractBeanFactory.java:315) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.AbstractBeanFactory.getBean(AbstractBeanFactory.java:199) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.DefaultListableBeanFactory.preInstantiateSingletons(DefaultListableBeanFactory.java:759) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.context.support.AbstractApplicationContext.finishBeanFactoryInitialization(AbstractApplicationContext.java:869) ~[spring-context-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.context.support.AbstractApplicationContext.refresh(AbstractApplicationContext.java:550) ~[spring-context-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.boot.SpringApplication.refresh(SpringApplication.java:762) ~[spring-boot-2.0.4.RELEASE.jar:2.0.4.RELEASE]
    at org.springframework.boot.SpringApplication.refreshContext(SpringApplication.java:398) ~[spring-boot-2.0.4.RELEASE.jar:2.0.4.RELEASE]
    at org.springframework.boot.SpringApplication.run(SpringApplication.java:330) ~[spring-boot-2.0.4.RELEASE.jar:2.0.4.RELEASE]
    at org.springframework.boot.builder.SpringApplicationBuilder.run(SpringApplicationBuilder.java:137) [spring-boot-2.0.4.RELEASE.jar:2.0.4.RELEASE]
    at com.integration.datastream.DatastreamApplication.main(DatastreamApplication.java:37) [main/:na]
Caused by: org.springframework.beans.factory.BeanCreationException: Error creating bean with name 'chunker': FactoryBean threw exception on object creation; nested exception is java.lang.IllegalArgumentException: targetObject must not be null
    at org.springframework.beans.factory.support.FactoryBeanRegistrySupport.doGetObjectFromFactoryBean(FactoryBeanRegistrySupport.java:178) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.FactoryBeanRegistrySupport.getObjectFromFactoryBean(FactoryBeanRegistrySupport.java:101) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.AbstractBeanFactory.getObjectForBeanInstance(AbstractBeanFactory.java:1645) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.AbstractAutowireCapableBeanFactory.getObjectForBeanInstance(AbstractAutowireCapableBeanFactory.java:1175) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.AbstractBeanFactory.doGetBean(AbstractBeanFactory.java:327) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.AbstractBeanFactory.getBean(AbstractBeanFactory.java:199) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.integration.config.annotation.AbstractMethodAnnotationPostProcessor.resolveTargetBeanFromMethodWithBeanAnnotation(AbstractMethodAnnotationPostProcessor.java:449) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.config.annotation.AbstractMethodAnnotationPostProcessor.postProcess(AbstractMethodAnnotationPostProcessor.java:133) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.config.annotation.MessagingAnnotationPostProcessor.processAnnotationTypeOnMethod(MessagingAnnotationPostProcessor.java:185) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.config.annotation.MessagingAnnotationPostProcessor.lambda$postProcessAfterInitialization$0(MessagingAnnotationPostProcessor.java:158) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.util.ReflectionUtils.doWithMethods(ReflectionUtils.java:562) ~[spring-core-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.util.ReflectionUtils.doWithMethods(ReflectionUtils.java:569) ~[spring-core-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.integration.config.annotation.MessagingAnnotationPostProcessor.postProcessAfterInitialization(MessagingAnnotationPostProcessor.java:139) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.beans.factory.support.AbstractAutowireCapableBeanFactory.applyBeanPostProcessorsAfterInitialization(AbstractAutowireCapableBeanFactory.java:431) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.AbstractAutowireCapableBeanFactory.initializeBean(AbstractAutowireCapableBeanFactory.java:1703) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.beans.factory.support.AbstractAutowireCapableBeanFactory.doCreateBean(AbstractAutowireCapableBeanFactory.java:573) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    ... 13 common frames omitted
Caused by: java.lang.IllegalArgumentException: targetObject must not be null
    at org.springframework.util.Assert.notNull(Assert.java:193) ~[spring-core-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    at org.springframework.integration.util.MessagingMethodInvokerHelper.<init>(MessagingMethodInvokerHelper.java:353) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.util.MessagingMethodInvokerHelper.<init>(MessagingMethodInvokerHelper.java:231) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.aggregator.MethodInvokingMessageListProcessor.<init>(MethodInvokingMessageListProcessor.java:60) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.aggregator.MethodInvokingMessageGroupProcessor.<init>(MethodInvokingMessageGroupProcessor.java:53) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.config.AggregatorFactoryBean.createHandler(AggregatorFactoryBean.java:176) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.config.AggregatorFactoryBean.createHandler(AggregatorFactoryBean.java:46) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.config.AbstractSimpleMessageHandlerFactoryBean.createHandlerInternal(AbstractSimpleMessageHandlerFactoryBean.java:185) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.config.AbstractSimpleMessageHandlerFactoryBean.getObject(AbstractSimpleMessageHandlerFactoryBean.java:173) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.integration.config.AbstractSimpleMessageHandlerFactoryBean.getObject(AbstractSimpleMessageHandlerFactoryBean.java:58) ~[spring-integration-core-5.0.7.RELEASE.jar:5.0.7.RELEASE]
    at org.springframework.beans.factory.support.FactoryBeanRegistrySupport.doGetObjectFromFactoryBean(FactoryBeanRegistrySupport.java:171) ~[spring-beans-5.0.8.RELEASE.jar:5.0.8.RELEASE]
    ... 28 common frames omitted

【问题讨论】:

    标签: spring spring-integration


    【解决方案1】:

    您的问题不清楚,这就是为什么人们已经对它投了反对票。为我们提供尽可能多的信息会很棒。例如,我相信您有某种filter 来检查行的长度,并且您还尝试使用aggregator 做一些事情。这就是您与提到的sequence size 斗争的地方。

    我认为你的想法是正确的,你仍然应该继续使用FileSplitter,但稍微改进了配置。首先,我们真的可能无法知道文件的原始大小。另一方面,您无法提前知道要过滤多少行。

    对于这样的任务,FileSplitter 建议使用markers 选项:https://docs.spring.io/spring-integration/docs/current/reference/html/files.html#file-splitter

    读取文件前的第一条消息是payload FileSplitter.FileMarkerMark.START。最后一个,当我们读完文件时是Mark.END。我建议在过滤线路长度之前有一个路由器 (PayloadTypeRouter),将路由器 String 有效负载放入线路长度过滤器,然后再到聚合器。那些FileSplitter.FileMarker 消息应该转到路由器的其他分支以过滤Mark.START 消息,因为我们对这种情况不感兴趣。 Mark.END 应该转到提到的聚合器,并且应该将这个确切地视为 RealeaseStrategy 以完成一组这些行。 FileSplitter.FileMarker 消息可以从文件分组功能中跳过,只发出你感兴趣的List&lt;String&gt;

    您可以在相关的 JIRA 中找到更多信息:https://jira.spring.io/browse/INT-4116

    更新

    由于您的任务略有不同,因此您一开始没有正确解释问题真是令人失望...

    对于当前的任务,我仍然会坚持使用FileSplitter

    @Splitter(inputChannel = "x")
    @Bean
    public FileSplitter fileSplitter() {
       FileSplitter fileSplitter = new FileSplitter();
       fileSplitter.setOutputChannelName("chunkingChannel");
       return fileSplitter;
    }
    

    并且将使用aggregator 和简单的MessageCountReleaseStrategy(1000)groupTimeout 作为最后一个块&lt; 1000

    @ServiceActivator(inputChannel = "chunkingChannel")
    @Bean
    public AggregatorFactoryBean chunker() {
        AggregatorFactoryBean aggregator = new AggregatorFactoryBean();
        aggregator.setReleaseStrategy(new MessageCountReleaseStrategy(1000));
        aggregator.setExpireGroupsUponCompletion(true);
        aggregator.setGroupTimeoutExpression(new ValueExpression<>(100L));
        aggregator.setOutputChannelName("y");
        aggregator.setProcessorBean(new DefaultAggregatingMessageGroupProcessor());
        return aggregator;
    }
    

    【讨论】:

    • 嘿Artem,首先感谢您的回答,但这不是我要找的东西,也许我的另一个问题让您回答这个问题,无论如何我已经修改了我的查询,看看您是否可以帮助我出去。谢谢
    • 请在我的回答中查看更新。
    • 您好 Artem,感谢您的快速回复。我使用了您的方法,但出现错误:bean name chunker 的 Bean creation error (Target Object cannot be null)。所以我使用 AggregateMessageHandler 而不是 AggregatorFactoryBean 并且错误得到解决并溢出工作。我面临的唯一问题是最后一个块的大小
    • 很高兴看到您自己解决了所有问题!是时候接受答案了。我还想查看有关错误的整个堆栈跟踪:也许我们需要在框架中修复一些东西
    • 以前,我尝试只访问一个目录。现在我正在尝试实现以下场景:我有一个触发器文件和一个保存在不同目录中的数据文件。只有当我收到一个触发文件时,我才应该能够访问数据文件,然后进行拆分和进一步的处理逻辑。此外,这种情况下会有一个触发器文件但有多个数据文件。所以拿到触发器文件后,我应该可以处理所有的数据文件了。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-05-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多