Đa xử lý và đa luồng trong Python

Oct 26 2020

Tôi có một chương trình python 1) Đọc từ một tệp rất lớn từ Đĩa (thời gian ~ 95%) và sau đó 2) Xử lý và cung cấp đầu ra tương đối nhỏ (thời gian ~ 5%). Chương trình này sẽ được chạy trên TeraByte tệp.

Bây giờ tôi đang tìm cách Tối ưu hóa Chương trình này bằng cách sử dụng Đa xử lý và Đa luồng. Nền tảng tôi đang chạy là Máy ảo với 4 Bộ xử lý trên Máy ảo.

Tôi dự định có một Quy trình lập lịch trình sẽ thực hiện 4 Quy trình (giống như các bộ xử lý) và sau đó Mỗi Quy trình nên có một số luồng vì hầu hết các phần là I / O. Mỗi Luồng sẽ xử lý 1 tệp và sẽ báo cáo kết quả cho Luồng chính và sau đó sẽ báo cáo lại cho Processer lên lịch qua IPC. Bộ lập lịch có thể xếp hàng đợi những thứ này và cuối cùng ghi chúng vào đĩa theo cách có thứ tự

Vì vậy, tự hỏi Làm thế nào một người quyết định số lượng Quy trình và Luồng để tạo cho kịch bản như vậy? Có một cách Toán học để tìm ra kết hợp tốt nhất là gì.

Cảm ơn bạn

Trả lời

6 Booboo Nov 04 2020 at 01:27

Tôi nghĩ tôi sẽ sắp xếp nó ngược lại với những gì bạn đang làm. Đó là, tôi sẽ tạo một nhóm luồng có kích thước nhất định sẽ chịu trách nhiệm tạo ra kết quả. Các tác vụ được gửi đến nhóm này sẽ được chuyển như một đối số mà nhóm bộ xử lý có thể được sử dụng bởi chuỗi công nhân để gửi các phần công việc bị ràng buộc bởi CPU. Nói cách khác, các nhân viên nhóm luồng chủ yếu sẽ thực hiện tất cả các hoạt động liên quan đến đĩa và giao cho nhóm bộ xử lý bất kỳ công việc nào đòi hỏi nhiều CPU.

Kích thước của nhóm bộ xử lý chỉ nên bằng số bộ xử lý bạn có trong môi trường của mình. Rất khó để đưa ra một kích thước chính xác cho nhóm chủ đề; nó phụ thuộc vào số lượng hoạt động đồng thời trên đĩa mà nó có thể xử lý trước khi quy luật trả về giảm dần có hiệu lực. Nó cũng phụ thuộc vào bộ nhớ của bạn: Pool càng lớn, tài nguyên bộ nhớ sẽ được sử dụng càng lớn, đặc biệt nếu toàn bộ tệp phải được đọc vào bộ nhớ để xử lý. Vì vậy, bạn có thể phải thử nghiệm với giá trị này. Đoạn mã dưới đây phác thảo những ý tưởng này. Những gì bạn thu được từ nhóm luồng là sự chồng chéo của các hoạt động I / O lớn hơn mức bạn sẽ đạt được nếu bạn chỉ sử dụng nhóm bộ xử lý nhỏ:

from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
from functools import partial
import os

def cpu_bound_function(arg1, arg2):
    ...
    return some_result



def io_bound_function(process_pool_executor, file_name):
    with open(file_name, 'r') as f:
        # Do disk related operations:
        . . . # code omitted
        # Now we have to do a CPU-intensive operation:
        future = process_pool_executor.submit(cpu_bound_function, arg1, arg2)
        result = future.result() # get result
        return result
    
file_list = [file_1, file_2, file_n]
N_FILES = len(file_list)
MAX_THREADS = 50 # depends on your configuration on how well the I/O can be overlapped
N_THREADS = min(N_FILES, MAX_THREADS) # no point in creating more threds than required
N_PROCESSES = os.cpu_count() # use the number of processors you have

with ThreadPoolExecutor(N_THREADS) as thread_pool_executor:
    with ProcessPoolExecutor(N_PROCESSES) as process_pool_executor:
        results = thread_pool_executor.map(partial(io_bound_function, process_pool_executor), file_list)

Lưu ý quan trọng :

Một cách tiếp cận khác đơn giản hơn nhiều là chỉ có một nhóm bộ xử lý duy nhất có kích thước lớn hơn số bộ xử lý CPU mà bạn có, ví dụ: 25. Các quy trình công nhân sẽ thực hiện cả hoạt động I / O và CPU. Mặc dù bạn có nhiều quy trình hơn CPU, nhưng nhiều quy trình sẽ ở trạng thái chờ đợi I / O hoàn tất cho phép chạy công việc sử dụng nhiều CPU.

Nhược điểm của phương pháp này là chi phí tạo N quá trình lớn hơn nhiều so với chi phí tạo N luồng + một số lượng nhỏ quy trình. Tuy nhiên, khi thời gian chạy của các nhiệm vụ được gửi đến nhóm ngày càng lớn hơn, thì chi phí tăng này sẽ giảm dần theo tỷ lệ phần trăm nhỏ hơn trong tổng thời gian chạy. Vì vậy, nếu nhiệm vụ của bạn không phải là tầm thường, đây có thể là một sự đơn giản hóa hợp lý.

Cập nhật: Điểm chuẩn của cả hai phương pháp

Tôi đã thực hiện một số điểm chuẩn so với hai phương pháp xử lý 24 tệp có kích thước xấp xỉ 10.000KB (thực tế, đây chỉ là 3 tệp khác nhau được xử lý 8 lần mỗi tệp, vì vậy có thể đã thực hiện một số bộ nhớ đệm):

Phương pháp 1 (nhóm luồng + nhóm bộ xử lý)

from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
from functools import partial
import os
from math import sqrt
import timeit


def cpu_bound_function(b):
    sum = 0.0
    for x in b:
        sum += sqrt(float(x))
    return sum

def io_bound_function(process_pool_executor, file_name):
    with open(file_name, 'rb') as f:
        b = f.read()
        future = process_pool_executor.submit(cpu_bound_function, b)
        result = future.result() # get result
        return result

def main():
    file_list = ['/download/httpd-2.4.16-win32-VC14.zip'] * 8 + ['/download/curlmanager-1.0.6-x64.exe'] * 8 + ['/download/Element_v2.8.0_UserManual_RevA.pdf'] * 8
    N_FILES = len(file_list)
    MAX_THREADS = 50 # depends on your configuration on how well the I/O can be overlapped
    N_THREADS = min(N_FILES, MAX_THREADS) # no point in creating more threds than required
    N_PROCESSES = os.cpu_count() # use the number of processors you have

    with ThreadPoolExecutor(N_THREADS) as thread_pool_executor:
        with ProcessPoolExecutor(N_PROCESSES) as process_pool_executor:
            results = list(thread_pool_executor.map(partial(io_bound_function, process_pool_executor), file_list))
            print(results)

if __name__ == '__main__':
    print(timeit.timeit(stmt='main()', number=1, globals=globals()))

Phương pháp 2 (chỉ nhóm bộ xử lý)

from concurrent.futures import ProcessPoolExecutor
from math import sqrt
import timeit


def cpu_bound_function(b):
    sum = 0.0
    for x in b:
        sum += sqrt(float(x))
    return sum

def io_bound_function(file_name):
    with open(file_name, 'rb') as f:
        b = f.read()
        result = cpu_bound_function(b)
        return result

def main():
    file_list = ['/download/httpd-2.4.16-win32-VC14.zip'] * 8 + ['/download/curlmanager-1.0.6-x64.exe'] * 8 + ['/download/Element_v2.8.0_UserManual_RevA.pdf'] * 8
    N_FILES = len(file_list)
    MAX_PROCESSES = 50 # depends on your configuration on how well the I/O can be overlapped
    N_PROCESSES = min(N_FILES, MAX_PROCESSES) # no point in creating more threds than required

    with ProcessPoolExecutor(N_PROCESSES) as process_pool_executor:
        results = list(process_pool_executor.map(io_bound_function, file_list))
        print(results)

if __name__ == '__main__':
    print(timeit.timeit(stmt='main()', number=1, globals=globals()))

Các kết quả:

(Tôi có 8 lõi)

Nhóm luồng + Nhóm xử lý: 13,5 giây Nhóm xử lý Một mình: 13,3 giây

Kết luận: Tôi sẽ thử cách tiếp cận đơn giản hơn trước tiên là chỉ sử dụng một nhóm bộ xử lý cho mọi thứ. Bây giờ, một chút khó khăn là quyết định số lượng quy trình tối đa cần tạo là bao nhiêu, đây là một phần của câu hỏi ban đầu của bạn và có một câu trả lời đơn giản khi tất cả những gì nó đang làm là tính toán nhiều CPU. Nếu số lượng tệp bạn đang đọc không quá nhiều, thì vấn đề cần lưu ý là; bạn có thể có một quy trình cho mỗi tệp. Nhưng nếu bạn có hàng trăm tệp, bạn sẽ không muốn có hàng trăm quy trình trong nhóm của mình (cũng có giới hạn trên về số lượng quy trình bạn có thể tạo và một lần nữa có những ràng buộc khó chịu về bộ nhớ). Không có cách nào tôi có thể cung cấp cho bạn một con số chính xác. Nếu bạn có một số lượng lớn tệp, hãy bắt đầu với kích thước nhóm nhỏ và tiếp tục tăng lên cho đến khi bạn không nhận được thêm lợi ích nào (tất nhiên, bạn có thể không muốn xử lý nhiều tệp hơn số lượng tối đa cho các thử nghiệm này hoặc bạn sẽ chạy mãi mãi chỉ cần quyết định kích thước hồ bơi tốt cho lần chạy thực sự).

Stack Oct 31 2020 at 15:53

Đối với xử lý song song: Tôi đã thấy câu hỏi này và trích dẫn câu trả lời được chấp nhận:

Trong thực tế, có thể rất khó để tìm ra số luồng tối ưu và thậm chí con số đó có thể sẽ thay đổi mỗi khi bạn chạy chương trình. Vì vậy, về mặt lý thuyết, số luồng tối ưu sẽ là số lõi bạn có trên máy của mình. Nếu lõi của bạn là "siêu phân luồng" (như Intel gọi nó), nó có thể chạy 2 luồng trên mỗi lõi. Sau đó, trong trường hợp đó, số luồng tối ưu sẽ gấp đôi số lõi trên máy của bạn.

Đối với đa xử lý: Ai đó đã hỏi một câu hỏi tương tự ở đây và câu trả lời được chấp nhận cho biết điều này:

Nếu tất cả các luồng / quy trình của bạn thực sự bị ràng buộc bởi CPU, thì bạn nên chạy càng nhiều quy trình mà CPU báo cáo lõi. Do HyperThreading, mỗi lõi CPU vật lý có thể có nhiều lõi ảo. Gọi multiprocessing.cpu_countđể lấy số lõi ảo.

Nếu chỉ p trong số 1 luồng của bạn bị ràng buộc bởi CPU, bạn có thể điều chỉnh số đó bằng cách nhân với p. Ví dụ: nếu một nửa quy trình của bạn bị ràng buộc bởi CPU (p = 0.5) và bạn có hai CPU với 4 lõi mỗi CPU và 2x HyperThreading, bạn nên bắt đầu quy trình 0,5 * 2 * 4 * 2 = 8.

Mấu chốt ở đây là hiểu bạn đang sử dụng máy gì, từ đó bạn có thể chọn một số luồng / quy trình gần như tối ưu để phân chia việc thực thi mã của bạn. Và tôi đã nói gần như tối ưu bởi vì nó sẽ thay đổi một chút mỗi khi bạn chạy tập lệnh của mình, vì vậy sẽ rất khó để dự đoán con số tối ưu này theo quan điểm toán học.

Đối với tình huống cụ thể của bạn, nếu máy của bạn có 4 lõi, tôi khuyên bạn chỉ nên tạo tối đa 4 luồng và sau đó chia chúng:

  • 1 đến chủ đề chính.
  • 3 để đọc và xử lý tệp.
Wilson.F Oct 31 2020 at 17:00

sử dụng nhiều quy trình để tăng tốc hiệu suất IO có thể không phải là một ý kiến ​​hay, hãy kiểm tra điều này và mã mẫu bên dưới nó để xem nó hữu ích hơn

FrançoisB. Nov 03 2020 at 04:51

Một ý tưởng có thể là có một luồng chỉ đọc tệp (Nếu tôi hiểu rõ, chỉ có một tệp) và đẩy các phần độc lập (ví dụ: các hàng) vào hàng đợi với các thông báo.
Các tin nhắn có thể được xử lý bởi 4 chuỗi. Bằng cách này, bạn có thể tối ưu hóa tải giữa các bộ xử lý.

vpelletier Nov 05 2020 at 06:54

Trên quy trình có ràng buộc I / O mạnh mẽ (như những gì bạn đang mô tả), bạn không nhất thiết phải đa luồng cũng như đa xử lý: bạn cũng có thể sử dụng các nguyên tắc I / O nâng cao hơn từ hệ điều hành của mình.

Ví dụ trên Linux, bạn có thể gửi yêu cầu đọc tới hạt nhân cùng với một bộ đệm có thể thay đổi có kích thước phù hợp và được thông báo khi bộ đệm được lấp đầy. Điều này có thể được thực hiện bằng cách sử dụng API AIO , mà tôi đã viết liên kết thuần python: python-libaio ( libaio trên pypi)) hoặc với API io_uring gần đây hơn mà dường như có liên kết python CFFI ( liburing trên pypy) (Tôi chưa sử dụng io_uring cũng như liên kết python này).

Điều này loại bỏ sự phức tạp của quá trình xử lý song song ở cấp độ của bạn, có thể giảm số lần chuyển đổi ngữ cảnh OS / vùng người dùng (giảm thời gian cpu hơn nữa) và cho phép hệ điều hành biết thêm về những gì bạn đang cố gắng thực hiện, tạo cơ hội cho việc lập lịch IO hiệu quả hơn (trong môi trường ảo hóa, tôi sẽ không ngạc nhiên nếu nó làm giảm số lượng bản sao dữ liệu, mặc dù tôi chưa tự mình thử).

Tất nhiên, nhược điểm là chương trình của bạn sẽ bị ràng buộc chặt chẽ hơn với hệ điều hành mà bạn đang thực thi nó, đòi hỏi nhiều nỗ lực hơn để làm cho nó chạy trên một hệ điều hành khác.