代码之家  ›  专栏  ›  技术社区  ›  Makky

TaskExecutor未在Spring集成中工作

  •  3
  • Makky  · 技术社区  · 7 年前

    我已经安装了带有任务执行器的文件轮询器

    ExecutorService executorService = Executors.newFixedThreadPool(10);
    
                LOG.info("Setting up the poller for directory {} ", finalDirectory);
                StandardIntegrationFlow standardIntegrationFlow = IntegrationFlows.from(new CustomFileReadingSource(finalDirectory),
                        c -> c.poller(Pollers.fixedDelay(5, TimeUnit.SECONDS, 5)
                                .taskExecutor(executorService)
                                .maxMessagesPerPoll(10)
                                .advice(new LoggerSourceAdvisor(finalDirectory))
                        ))
    
    
                        //move file to processing first processing                    
                        .transform(new FileMoveTransformer("C:/processing", true))
                        .channel("fileRouter")
                        .get();
    

    threadpool 每次轮询最多10条消息,最多10条消息。如果我放10个文件,它仍然会一个接一个地处理。这里可能出了什么问题?

    *更新*

    在Gary的回答之后,它工作得非常好,尽管我现在还有其他问题。

    setDirectory(new File(path));
            DefaultDirectoryScanner scanner = new DefaultDirectoryScanner();
    
            scanner.setFilter(new AcceptAllFileListFilter<>());
            setScanner(scanner);
    

    使用的原因 AcceptAll 因为同一个文件可能会再次出现,所以我会先移动文件。但是当我启用线程执行器时,同一个文件正在由多个线程处理,我假设是因为 AcceptAllFile

    如果我换成 AcceptOnceFileListFilter

    问题/缺陷

    课堂上 AbstractPersistentAcceptOnceFileListFilter 我们有这个密码

    @Override
        public boolean accept(F file) {
            String key = buildKey(file);
            synchronized (this.monitor) {
                String newValue = value(file);
                String oldValue = this.store.putIfAbsent(key, newValue);
                if (oldValue == null) { // not in store
                    flushIfNeeded();
                    return true;
                }
                // same value in store
                if (!isEqual(file, oldValue) && this.store.replace(key, oldValue, newValue)) {
                    flushIfNeeded();
                    return true;
                }
                return false;
            }
        }
    

    现在,例如,如果我设置了max per poll 5,并且有两个文件,那么两个线程可能会拾取相同的文件。

    但是另一条线到达了 accept

    如果文件不存在,则它将返回lastModified time为0,并返回true。

    这会导致问题,因为文件不存在。

    如果为0,则应返回false,因为该文件不再存在。

    1 回复  |  直到 7 年前
        1
  •  5
  •   Community Mohan Dere    6 年前

    将任务执行器添加到轮询器时;所做的只是调度程序线程将轮询任务交给线程池中的一个线程;这个 maxMessagesPerPoll 是轮询任务的一部分。轮询器本身每5秒只运行一次。要获得您想要的,您应该向流中添加一个执行器通道。。。

    @SpringBootApplication
    public class So53521593Application {
    
        private static final Logger logger = LoggerFactory.getLogger(So53521593Application.class);
    
        public static void main(String[] args) {
            SpringApplication.run(So53521593Application.class, args);
        }
    
        @Bean
        public IntegrationFlow flow() {
            ExecutorService exec = Executors.newFixedThreadPool(10);
            return IntegrationFlows.from(() -> "foo", e -> e
                        .poller(Pollers.fixedDelay(5, TimeUnit.SECONDS)
                                .maxMessagesPerPoll(10)))
                    .channel(MessageChannels.executor(exec))
                    .<String>handle((p, h) -> {
                        try {
                            logger.info(p);
                            Thread.sleep(10_000);
                        }
                        catch (InterruptedException e1) {
                            Thread.currentThread().interrupt();
                        }
                        return null;
                    })
                    .get();
        }
    }
    

    编辑

    这对我来说很好。。。

    @Bean
    public IntegrationFlow flow() {
        ExecutorService exec = Executors.newFixedThreadPool(10);
        return IntegrationFlows.from(Files.inboundAdapter(new File("/tmp/foo")).filter(
                    new FileSystemPersistentAcceptOnceFileListFilter(new SimpleMetadataStore(), "foo")),
                        e -> e.poller(Pollers.fixedDelay(5, TimeUnit.SECONDS)
                                .maxMessagesPerPoll(10)))
                .channel(MessageChannels.executor(exec))
                .handle((p, h) -> {
                    try {
                        logger.info(p.toString());
                        Thread.sleep(10_000);
                    }
                    catch (InterruptedException e1) {
                        Thread.currentThread().interrupt();
                    }
                    return null;
                })
                .get();
    }
    

    2018-11-28 11:46:05.196信息57607---[pool-1-thread-1]com.example.so53521593应用程序:/tmp/foo/test1.txt

    和 touch test1.txt

    2018-11-28 11:48:00.284信息57607---[pool-1-thread-3]com.example.so53521593应用程序:/tmp/foo/test1.txt

    编辑1

    @Bean
    public IntegrationFlow flow() {
        ExecutorService exec = Executors.newFixedThreadPool(10);
        return IntegrationFlows.from(Files.inboundAdapter(new File("/tmp/foo")).filter(
                    new FileSystemPersistentAcceptOnceFileListFilter(new SimpleMetadataStore(), "foo")),
                        e -> e.poller(Pollers.fixedDelay(5, TimeUnit.SECONDS)
                                .maxMessagesPerPoll(10)))
                .channel(MessageChannels.executor(exec))
                .<File>handle((p, h) -> {
                    try {
                        p.delete();
                        logger.info(p.toString());
                        Thread.sleep(10_000);
                    }
                    catch (InterruptedException e1) {
                        Thread.currentThread().interrupt();
                    }
                    return null;
                })
                .get();
    }
    

    和

    2018-11-28 13:22:23.690信息75681---[pool-1-thread-2]com.example.so53521593应用程序:/tmp/foo/test2.txt

    2018-11-28 13:22:23.690信息75681---[pool-1-thread-3]com.example.so53521593应用程序:/tmp/foo/test1.txt