Spring Integration File reading

2019-05-18 19:53发布

I am newbie to Spring Integration. I am working on solution, but I am stuck on a specific issue while using inbound file adapter ( FileReadingMessageSource ). I have to read files from different directories and process them and save the files in different directories. As I understand, the directory name is fixed at the start of the flow. Can some one help me on changing the directory name for different requests.

I attempted the following. First of all, I am not sure whether it is correct way to about and although it worked for only one directory. I think Poller was waiting for more files and never came back to read another directory.

@SpringBootApplication
@EnableIntegration
@IntegrationComponentScan
public class SiSampleFileProcessor {

    @Autowired
    MyFileProcessor myFileProcessor;

    @Value("${si.outdir}")
    String outDir;

    @Autowired
    Environment env;

    public static void main(String[] args) throws IOException {
        ConfigurableApplicationContext ctx = new SpringApplication(SiSampleFileProcessor.class).run(args);
        FileProcessingService gateway = ctx.getBean(FileProcessingService.class);
        boolean process = true;
        while (process) {
            System.out.println("Please enter the input Directory: ");
            String inDir = new Scanner(System.in).nextLine();
            if ( inDir.isEmpty() || inDir.equals("exit") ) {
                process=false;
            } else {
                System.out.println("Processing... " + inDir);
                gateway.processFilesin(inDir);
            }
        }
        ctx.close();
    }

    @MessagingGateway(defaultRequestChannel="requestChannel")
    public interface FileProcessingService {
        String processFilesin( String inputDir );
    }

    @Bean(name = PollerMetadata.DEFAULT_POLLER)
    public PollerMetadata poller() {                                      
    return Pollers.fixedDelay(1000).get();
    }

    @Bean
    public MessageChannel requestChannel() {
        return new DirectChannel();
    }

    @ServiceActivator(inputChannel = "requestChannel")
    @Bean 
    GenericHandler<String> fileReader() {
        return new GenericHandler<String>() {
            @Override
            public Object handle(String p, Map<String, Object> map) {
                FileReadingMessageSource fileSource = new FileReadingMessageSource();
                fileSource.setDirectory(new File(p));
                Message<File> msg;
                while( (msg = fileSource.receive()) != null ) {
                    fileInChannel().send(msg);
                }
                return null;  // Not sure what to return!
            }
        };
    }

    @Bean
    public MessageChannel fileInChannel() {
        return MessageChannels.queue("fileIn").get();
    }

    @Bean
    public IntegrationFlow fileProcessingFlow() {
        return IntegrationFlows.from(fileInChannel())
                .handle(myFileProcessor)
                .handle(Files.outboundAdapter(new File(outDir)).autoCreateDirectory(true).get())
                .get();
    }    
}

EDIT: Based on Gary's response replaced some methods as

@MessagingGateway(defaultRequestChannel="requestChannel")
public interface FileProcessingService {
    boolean processFilesin( String inputDir );
}

@ServiceActivator(inputChannel = "requestChannel")
public boolean fileReader(String inDir) {
    FileReadingMessageSource fileSource = new FileReadingMessageSource();
    fileSource.setDirectory(new File(inDir));
    fileSource.afterPropertiesSet();
    fileSource.start();
    Message<File> msg;
    while ((msg = fileSource.receive()) != null) {
        fileInChannel().send(msg);
    }
    fileSource.stop();
    System.out.println("Sent all files in directory: " + inDir);
    return true;
}

Now it is working as expected.

2条回答
小情绪 Triste *
2楼-- · 2019-05-18 20:32

You can use this code

FileProcessor.java

import org.springframework.messaging.Message;
import org.springframework.stereotype.Component;
@Component
public class FileProcessor {

    private static final String HEADER_FILE_NAME = "file_name";
    private static final String MSG = "%s received. Content: %s";

    public void process(Message<String> msg) {
        String fileName = (String) msg.getHeaders().get(HEADER_FILE_NAME);
        String content = msg.getPayload();
        //System.out.println(String.format(MSG, fileName, content));
        System.out.println(content);

    }
}

LastModifiedFileFilter.java

package com.example.demo;

import org.springframework.integration.file.filters.AbstractFileListFilter;

import java.io.File;
import java.util.HashMap;
import java.util.Map;

public class LastModifiedFileFilter extends AbstractFileListFilter<File> {
    private final Map<String, Long> files = new HashMap<>();
    private final Object monitor = new Object();

    @Override
    protected boolean accept(File file) {
        synchronized (this.monitor) {
            Long previousModifiedTime = files.put(file.getName(), file.lastModified());

            return previousModifiedTime == null || previousModifiedTime != file.lastModified();
        }
    }
}

Main Class= DemoApplication.java

package com.example.demo;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration;
import org.springframework.boot.autoconfigure.orm.jpa.HibernateJpaAutoConfiguration;
import java.io.File;
import java.io.IOException;
import java.nio.charset.Charset;
import org.apache.commons.io.FileUtils;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.annotation.Aggregator;
import org.springframework.integration.annotation.InboundChannelAdapter;
import org.springframework.integration.annotation.Poller;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.core.MessageSource;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.integration.dsl.IntegrationFlows;
import org.springframework.integration.dsl.channel.MessageChannels;
import org.springframework.integration.dsl.core.Pollers;
import org.springframework.integration.file.FileReadingMessageSource;
import org.springframework.integration.file.filters.CompositeFileListFilter;
import org.springframework.integration.file.filters.SimplePatternFileListFilter;
import org.springframework.integration.file.transformer.FileToStringTransformer;
import org.springframework.integration.scheduling.PollerMetadata;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.PollableChannel;
import org.springframework.stereotype.Component;


@SpringBootApplication
@Configuration
public class DemoApplication {

    private static final String DIRECTORY = "E:/usmandata/logs/input/";

    public static void main(String[] args) throws IOException, InterruptedException {
        SpringApplication.run(DemoApplication.class, args);


    }


    @Bean
    public IntegrationFlow processFileFlow() {
        return IntegrationFlows
                .from("fileInputChannel")
                .transform(fileToStringTransformer())
                .handle("fileProcessor", "process").get();
    }

    @Bean
    public MessageChannel fileInputChannel() {
        return new DirectChannel();
    }

    @Bean
    @InboundChannelAdapter(value = "fileInputChannel", poller = @Poller(fixedDelay = "1000"))
    public MessageSource<File> fileReadingMessageSource() {
        CompositeFileListFilter<File> filters =new CompositeFileListFilter<>();
        filters.addFilter(new SimplePatternFileListFilter("*.log"));
        filters.addFilter(new LastModifiedFileFilter());

        FileReadingMessageSource source = new FileReadingMessageSource();
        source.setAutoCreateDirectory(true);
        source.setDirectory(new File(DIRECTORY));
        source.setFilter(filters);

        return source;
    }

    @Bean
    public FileToStringTransformer fileToStringTransformer() {
        return new FileToStringTransformer();
    }

    @Bean
    public FileProcessor fileProcessor() {
        return new FileProcessor();
    }
}
查看更多
不美不萌又怎样
3楼-- · 2019-05-18 20:40

The FileReadingMessageSource uses a DirectoryScanner internally; it is normally set up by Spring after the properties are injected. Since you are managing the object outside of Spring, you need to call Spring bean initialization and lifecycle methods afterPropertiesSet() , start() and stop().

Call stop() when the receive returns null.

> return null;  // Not sure what to return!

If you return nothing, your calling thread will hang in the gateway waiting for a response. You could change the gateway to return void or, since your gateway is expecting a String, just return some value.

However, your calling code is not looking at the result anyway.

> gateway.processFilesin(inDir);

Also, remove the @Bean from the @ServiceActivator; with that style, the bean type must be MessageHandler.

查看更多
登录 后发表回答