【问题标题】:spring boot API - document processing and executing python script on documents in parallelspring boot API - 文档处理并在文档上并行执行python脚本
【发布时间】:2020-02-15 18:29:50
【问题描述】:

场景

  1. 在我的应用程序中,有 3 个进程正在将共享驱动器上的文档复制到各自文件夹中。
  2. 一旦任何文档(通过任何进程)复制到共享驱动器上,目录观察程序 (Java) 代码就会获取该文档并使用“进程”调用 Python 脚本并对文档进行一些处理。代码sn-p如下:

    Process pr = Runtime.getRuntime().exec(pythonCommand);
                // retrieve output from python script
                BufferedReader bfr = new BufferedReader(new InputStreamReader(pr.getInputStream()));
                String line = "";
                while ((line = bfr.readLine()) != null) {
                    // display each output line from python script
                    logger.info(line);
                }
                pr.waitFor();
    
  3. 目前我的代码等待文档上的 python 代码执行完成。只有在那之后它才会拿起下一个文件。 Python 代码需要 30 秒才能完成

  4. 处理文档后,文档从当前文件夹移动到存档或错误文件夹。
  5. 请在下面找到该场景的屏幕截图:

有什么问题?

  1. 我的代码正在按顺序处理文档,我需要并行处理文档。
  2. 由于 Python 代码大约需要 30 秒,因此目录观察程序创建的一些事件也会丢失。
  3. 如果在短时间内收到大约 400 个文档,文档处理将停止。

我在寻找什么?

  1. 设计用于并行处理文档的解决方案。
  2. 如果文档处理出现任何故障情况,必须自动处理待处理的文档。
  3. 我也尝试了 spring boot schedular,但仍然只能按顺序处理文档。
  4. 是否可以将 Python 代码作为后台进程并行调用。

很抱歉这个问题太长了,但我已经被困在这个问题上很多天了,并且已经看过很多类似的问题。 谢谢!

【问题讨论】:

    标签: java python spring-boot document


    【解决方案1】:

    一种选择是使用 JDK 提供的ExecutorService,它可以执行RunnableCallable 任务。您需要创建一个实现Runnable 的类,该类将执行您的Python 脚本,在收到新文档后,您需要创建该类的新实例并将其传递给ExecutorService

    为了展示其工作原理,我们将使用一个简单的 Python 脚本,该脚本将线程名称作为参数,打印其执行的开始时间,休眠 10 秒并打印结束时间:

    import time
    import sys
    
    print "%s start : %s" % (sys.argv[1], time.ctime())
    time.sleep(10)
    print "%s end : %s" % (sys.argv[1], time.ctime())
    

    首先,我们实现运行脚本的类,并将构造函数中获得的名称传递给它:

    class ScriptRunner implements Runnable {
    
        private String thread;
    
        ScriptRunner(String thread) {
            this.thread = thread;
        }
    
        @Override
        public void run() {
            try {
                ProcessBuilder ps = new ProcessBuilder("py", "test.py", thread);
                ps.redirectErrorStream(true);
                Process pr = ps.start();
                try (BufferedReader in = new BufferedReader(new InputStreamReader(pr.getInputStream()))) {
                    String line;
                    while ((line = in.readLine()) != null) {
                        System.out.println(line);
                    }
                }
                pr.waitFor();
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
    }
    

    然后我们创建main 方法,该方法使用固定数量的5 个并行线程创建ExecutorService,并将10 个ScriptRunner 实例传递给它,中断时间为1 秒:

    public static void main(String[] args) throws InterruptedException {
        ExecutorService executor = Executors.newFixedThreadPool(5);
        for (int i = 1; i <= 10; i++) {
            executor.submit(new ScriptRunner("Thread_" + i));
            Thread.sleep(1000);
        }
        executor.shutdown();
    }
    

    如果我们运行这个方法,我们会看到服务,由于指定的限制,最多有5个并行运行的任务,其余的都进入队列并在释放的线程中启动:

    Thread_1 start : Sat Nov 23 11:40:14 2019
    Thread_1 end : Sat Nov 23 11:40:24 2019    // the first task is completed..
    Thread_2 start : Sat Nov 23 11:40:15 2019
    ...
    Thread_5 end : Sat Nov 23 11:40:28 2019
    Thread_6 start : Sat Nov 23 11:40:24 2019  // ..and the sixth is started
    ...
    Thread_10 end : Sat Nov 23 11:40:38 2019
    

    【讨论】:

    • 非常感谢您的帮助!最初,我遇到了错误,但最终我能够解决该错误。这对我有用。
    • 仅供参考 - 我只在 Python 脚本中遇到错误,在 PyCharm 中,“%s”未被识别。所以我使用了下面的脚本,之后这个工作: import datetime import time startTime = datetime.datetime.now() print("start =", startTime) time.sleep(10) endTime = datetime.datetime.now() print( "end ===", endTime)
    【解决方案2】:
    • 创建两个队列(阻塞队列):
      执行队列
      错误队列

    • 创建两个线程(您可以根据需要创建任意多个):
      第一线程
      第二线程

    • 概念:
      生产者-消费者

    • 生产者(目录观察线程):
      目录观察者

    • 消费者:
      第一线程
      第二线程

    详情:

    • 两个队列的增删方法必须同步。某一时刻只有一个线程会访问该方法。如果一个线程正在访问关键区域(生产者或消费者),其余线程将等待轮到它们。
    • 第一个生产者将开始工作,最初,消费者处于睡眠阶段。

      为什么?同步运行整个系统。

      你将如何得到它?处理后的睡眠生产者线程以及在作业开始时消费者睡眠的情况。

    • 第一个生产者或消费者将获取队列中的锁,处理工作并释放它。在这期间,如果任何线程(生产者或消费者)来获取数据,它们将等待轮到它们(使用线程池的概念)。

    • 只要将任何文档复制到共享驱动器(由任何进程),目录观察者(生产者)代码就会选择该文档的路径并同步存储在 executionQueue 中。

    • 现在Consumer会来取数据,FirstThread首先醒来,去executionQueue取数据。 FirstThread 将获取锁定的 executionQueue,然后获取数据并释放其中的锁。如果在 SecondThread 之间来取数据,它将等待轮到他。

    • 从 executionQueue 中获取数据后,FirstThread 将从 location 中获取文档并使用获取的文档调用 Python 脚本。

    • 在 SecondThread 之间会获取锁并获取路径并开始处理与 FirstThread 相同的概念。

    • 几秒钟后 FirstThread 将完成他的工作,然后它将进入 executionQueue 并再次获取锁并获取文件路径并释放锁并开始处理相同的工作并为 SecondThread 休息太……

    • 在处理该文件时,如果发生任何错误,则将该路径信息发送到 errorQueue 方法,并在一天结束时或系统空闲时使用相同的概念或手动分析该 errorQueue 信息。

    • 如果 executionQueue 中没有可用数据,则此时生产者线程(Directory watcher)已经处于睡眠阶段。然后消费者线程会来到executionQueue取数据,他们不会得到任何数据并进入睡眠阶段,如1分钟,1分钟后它会再次醒来并去取数据等等......

    • 在每个步骤日志中,信息将帮助您更好地进行分析。

    • 使用该概念,您可以并行运行整个系统。

    【讨论】:

    • 感谢您为解决问题所做的努力!
    【解决方案3】:

    你可以试试python中的多处理模块here

    由于 GIL,Python 的线程不会加速计算 受 CPU 限制。

    这个问题Solving embarassingly parallel problems using Python multiprocessing的可能重复

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-12-12
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-10-18
      • 2021-10-31
      相关资源
      最近更新 更多