【问题标题】:Which is the suitable approach to handle multiple jobs in parallel without blocking the main thread这是在不阻塞主线程的情况下并行处理多个作业的合适方法
【发布时间】:2018-04-18 02:43:43
【问题描述】:

我有一个要求,即单个进程应该并行处理多个作业。这些作业中的每一个都应定期运行(例如,每 10 秒)。主线程也需要观察停止信号,当收到时应该停止所有线程并退出。

以下是我处理此要求的方法。

主线程将为每个作业创建线程并等待停止信号。每个线程负责处理循环机制。当主线程接收到停止信号时,它会向线程发送信号停止。

在这种机制下,踏板将一直运行。所以我在想可能有更好的方法来处理这个问题。 也许就像,主线程将跟踪每个作业应该何时执行,并在需要时启动线程。线程将执行一些操作并在完成时退出。在这种情况下,线程不会一直运行。

我目前不知道如何实施替代方法。所以,我想知道以上哪一个可能是好方法?另外,如果有其他更好的选择,请提出建议。

编辑 [3 月 19 日]:

也许我一开始就应该提到这一点,但这里是这样。

假设如果我有 2 个作业,则不需要同时运行这两个作业。例如,作业 1 应每 10 秒运行一次,作业 2 应每 20 秒运行一次。

此外,如果 Job 本身需要更多时间,则必须有某种机制来正确识别和处理它。也许跳过执行或等待上一个作业完成然后重新开始。目前,这一项的要求并不明确。但是,我应该能够识别和处理这种情况。

【问题讨论】:

标签: c++ multithreading timer


【解决方案1】:

这里有一些建议。

首先让主线程启动线程以根据需要完成任务。主线程在创建线程时将其分离,然后完全忘记它。分离的线程在返回时会自行清理。这可以很容易地简化主线程,并消除跟踪和管理不同线程的开销。如果线程几乎完全独立并且永远不会遇到无限循环,则此方法有效。主线程将在收到停止信号时退出,进而终止进程,但这取决于操作系统,请查看What happens to a detached thread when main() exits?

您启动了一个充当任务/线程管理器的线程。该对象/线程将根据需要创建新线程,并通过一些池或线程跟踪正确地等待它们。然后主线程可以通过互斥锁保护标志轻松地向线程跟踪器发送消息以停止,在这种情况下,线程跟踪器只需等待其所有产生的线程然后死亡。这需要比上述解决方案更多的开销,但它为您提供了有关哪些线程何时运行的更多信息。此外,您还可以更好地控制是否需要直接终止线程或在必要时如何向线程发送消息。它还允许安全清理,这取决于您的操作系统,线程可能会因进程死亡而缩短。如果线程需要能够更轻松地接收消息,并且如果您想确保线程即使在停止信号之后也能运行完成(想想 R/W 操作),这会更好。

您也可以将主线程和线程/管理器混合为一个,但这会使主线程更加复杂,并且对于大多数通用场景而言并不会显着减少开销。

【讨论】:

  • 感谢您的评论。在线程管理器的情况下(我正在考虑线程池,因为每次创建线程都非常昂贵),如果线程花费的时间比预期返回的时间长,有没有办法停止执行。
  • 你可以有一个定时器,可以设置一个标志(如共享布尔值),或者你可以让定时器直接杀死线程,但是要做到这一点,你需要使用本机线程句柄,即不同的系统不同,在 POSIX 上它的 pthread
  • 在杀死线程时我需要特别注意一些事情。我不知道杀死线程是否会导致任何不一致的行为。
  • 这取决于线程正在做什么以及在操作系统上,如果线程有文件打开不礼貌地杀死它可能会在进程的文件指针表中留下那些打开的文件。它还可能导致任何分配的堆内存泄漏,因为它不会被调用delete
【解决方案2】:

看看以下以给定间隔运行函数的类:

class CallbackTimer
{
public:
    ~CallbackTimer()
    {
        stop();
    }

    void start(std::chrono::milliseconds interval, std::function<void()> callback)
    {
        stop();
        shouldQuit = false;

        handle = std::async([=,callback=std::move(callback)]() {

            while (!shouldQuit)
            {
                auto nextStart = std::chrono::steady_clock::now() + interval;
                callback();
                std::this_thread::sleep_until(nextStart);
            }
        });
    }

    void stop()
    {
        if (handle.valid())
        {
            shouldQuit = true;
            handle.get();
        }
    }

private:
    std::atomic_bool shouldQuit;
    std::future<void> handle;
};

在 start() 上创建一个新线程,它以给定的时间间隔运行给定的“回调”。对于 10 秒的间隔,回调恰好每 10 秒调用一次,除非在下次启动到期时它仍在工作。在这种情况下,它会立即再次运行。

stop() 设置退出标志并等待线程退出。主路由的这个非常简单的实现不会中断 shouldQuit-checks 的睡眠,因此 CallbackTimer 退出可能需要一个完整的时间间隔。

请注意,回调时间可能会略有偏差,因此如果您以相同的时间间隔启动其中两个,则它们可能不会在一段时间后同时运行。但漂移应该是最小的,如果它很严重,你应该考虑另一种解决方案。

这是一个用法示例:

int main()
{
    CallbackTimer a;
    a.start(std::chrono::seconds(1), []() { std::cout << "a" << std::endl; });

    CallbackTimer b;
    b.start(std::chrono::seconds(2), []() { std::cout << "b" << std::endl; });

    std::this_thread::sleep_for(std::chrono::seconds(10));

    a.stop();
    b.stop();

    return 0;
}

为了更快的 stop(),您可以实现“beginStop()”方法,该方法仅将“shouldQuit”标志设置为 true,并在所有实例上调用“stop()”之前为所有实例调用此方法。这样第二个实例的关闭不会延迟到第一个实例的关闭完成。

【讨论】:

  • 非常感谢。我能够理解并提供考虑您的想法的解决方案。如果您能回答,我还有一个问题,如果线程内的 API 调用花费的时间比预期的多,是否有任何技术可以处理。目前我想停止/忽略它并再次开始一个新的请求。
  • 很抱歉,取消 API 调用没有“标准”方式。您可以发送一个新调用(使用不同的线程方法),但这可能会导致大量挂起的调用。但是,在 Windows-API 中有一个 TerminateThread 函数 (msdn.microsoft.com/en-us/library/windows/desktop/…) 可以取消线程。但由于目标线程无法清理和释放资源,因此几乎不可能安全使用它。
  • 我明白了。我也不想突然终止线程。但是,目前我不知道如何处理所有场景。无论如何,这是我担心的。谢谢你的回答。
【解决方案3】:

如果您的主线程只需要管理线程,您可以每 10 秒启动一次新线程并等待它们完成。

以下示例并行运行两个线程,等待它们完成并在循环之前再等待一秒。 (shouldQuit() 函数使用 Visual C++ 中的编译器特定函数;也许您必须在其中插入自己的代码)

#include <future>
#include <iostream>
#include <conio.h>
#include <thread>
#include <chrono>

int a()
{
    std::cout << 'a' << std::flush;
    return 1;
}

int b()
{
    std::cout << 'b' << std::flush;
    return 2;
}

bool shouldQuit()
{
    while (_kbhit())
    {
        if (_getch() == 27)
            return true;
    }
    return false;
}

int main()
{
    auto threadFunctions = { a, b };

    while(!shouldQuit())
    {
        auto threadFutures = std::vector<std::future<int>>{};

        // Run asynchronous tasks
        for (auto& threadFunction: threadFunctions)
            threadFutures.push_back(std::async(threadFunction));

        // Wait for all tasks to complete
        for (auto& threadFuture : threadFutures)
            threadFuture.get();

        // Wait a second
        std::this_thread::sleep_for(std::chrono::seconds(1));
    }

    return 0;
}

线程消耗的时间(当我们等待它们完成时)不是计时的一部分。对于运行时间最长的线程,线程调用之间的暂停时间为 1 秒,对于所有其他线程,至少为 1 秒。

注意:每隔几秒运行一个新线程并不是最有效的处理方式。在大多数操作系统上,创建新线程是一项昂贵的操作。如果性能至关重要,您应该按照问题中的描述重用现有线程。但是程序的整体复杂性会增加。

使用std::asyncstd::future 的好处是可以正确处理异常。如果一个线程抛出一个异常,这个异常会被捕获并存储在 std::future 中。它在调用future.get() 的线程中被重新抛出。因此,如果您的线程可能抛出异常,最好将您的 future.get() 调用包装在 try .. catch 中。

future.get() 调用将调用线程与创建未来的线程同步。如果线程已经完成,它只返回返回值。如果线程仍在工作,则调用线程将被阻塞,直到线程完成。

【讨论】:

  • 谢谢@Andreas H。在决定我要选择的方法时,我会考虑你的建议。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-11-08
  • 1970-01-01
  • 2014-05-25
  • 1970-01-01
相关资源
最近更新 更多