介紹
Python的multiprocessing模塊不但支持多進程,其中managers子模塊還支持把多進程分布到多臺機器上。一個服務進程可以作為調度者,將任務分布到其他多個機器的多個進程中,依靠網絡通信。
想到這,就在想是不是可以使用此模塊來實現一個簡單的作業調度系統。
實現
Job
首先創建一個Job類,為了測試簡單,只包含一個job id屬性
job.py
#!/usr/bin/env python # -*- coding: utf-8 -*- class Job: def __init__(self, job_id): self.job_id = job_id
Master
Master用來派發作業和顯示運行完成的作業信息
master.py
#!/usr/bin/env python # -*- coding: utf-8 -*- from Queue import Queue from multiprocessing.managers import BaseManager from job import Job
class Master:
def __init__(self): # 派發出去的作業隊列 self.dispatched_job_queue = Queue() # 完成的作業隊列 self.finished_job_queue = Queue() def get_dispatched_job_queue(self): return self.dispatched_job_queue def get_finished_job_queue(self): return self.finished_job_queue def start(self): # 把派發作業隊列和完成作業隊列注冊到網絡上 BaseManager.register('get_dispatched_job_queue', callable=self.get_dispatched_job_queue) BaseManager.register('get_finished_job_queue', callable=self.get_finished_job_queue) # 監聽端口和啟動服務 manager = BaseManager(address=('0.0.0.0', 8888), authkey='jobs') manager.start() # 使用上面注冊的方法獲取隊列 dispatched_jobs = manager.get_dispatched_job_queue() finished_jobs = manager.get_finished_job_queue() # 這里一次派發10個作業,等到10個作業都運行完后,繼續再派發10個作業 job_id = 0 while True: for i in range(0, 10): job_id = job_id + 1 job = Job(job_id) print('Dispatch job: %s' % job.job_id) dispatched_jobs.put(job) while not dispatched_jobs.empty(): job = finished_jobs.get(60) print('Finished Job: %s' % job.job_id) manager.shutdown() if __name__ == "__main__": master = Master() master.start()
Slave
Slave用來運行master派發的作業并將結果返回
slave.py
#!/usr/bin/env python # -*- coding: utf-8 -*- import time from Queue import Queue from multiprocessing.managers import BaseManager from job import Job
class Slave:
def __init__(self): # 派發出去的作業隊列 self.dispatched_job_queue = Queue() # 完成的作業隊列 self.finished_job_queue = Queue()
def start(self):
# 把派發作業隊列和完成作業隊列注冊到網絡上 BaseManager.register('get_dispatched_job_queue') BaseManager.register('get_finished_job_queue') # 連接master server = '127.0.0.1' print('Connect to server %s...' % server) manager = BaseManager(address=(server, 8888), authkey='jobs') manager.connect() # 使用上面注冊的方法獲取隊列 dispatched_jobs = manager.get_dispatched_job_queue() finished_jobs = manager.get_finished_job_queue() # 運行作業并返回結果,這里只是模擬作業運行,所以返回的是接收到的作業 while True: job = dispatched_jobs.get(timeout=1) print('Run job: %s ' % job.job_id) time.sleep(1) finished_jobs.put(job) if __name__ == "__main__": slave = Slave() slave.start()
測試
分別打開三個linux終端,第一個終端運行master,第二個和第三個終端用了運行slave,運行結果如下
master
$ python master.py Dispatch job: 1 Dispatch job: 2 Dispatch job: 3 Dispatch job: 4 Dispatch job: 5 Dispatch job: 6 Dispatch job: 7 Dispatch job: 8 Dispatch job: 9 Dispatch job: 10 Finished Job: 1 Finished Job: 2 Finished Job: 3 Finished Job: 4 Finished Job: 5 Finished Job: 6 Finished Job: 7 Finished Job: 8 Finished Job: 9 Dispatch job: 11 Dispatch job: 12 Dispatch job: 13 Dispatch job: 14 Dispatch job: 15 Dispatch job: 16 Dispatch job: 17 Dispatch job: 18 Dispatch job: 19 Dispatch job: 20 Finished Job: 10 Finished Job: 11 Finished Job: 12 Finished Job: 13 Finished Job: 14 Finished Job: 15 Finished Job: 16 Finished Job: 17 Finished Job: 18 Dispatch job: 21 Dispatch job: 22 Dispatch job: 23 Dispatch job: 24 Dispatch job: 25 Dispatch job: 26 Dispatch job: 27 Dispatch job: 28 Dispatch job: 29 Dispatch job: 30
slave1
$ python slave.py Connect to server 127.0.0.1... Run job: 1 Run job: 2 Run job: 3 Run job: 5 Run job: 7 Run job: 9 Run job: 11 Run job: 13 Run job: 15 Run job: 17 Run job: 19 Run job: 21 Run job: 23
slave2
$ python slave.py Connect to server 127.0.0.1... Run job: 4 Run job: 6 Run job: 8 Run job: 10 Run job: 12 Run job: 14 Run job: 16 Run job: 18 Run job: 20 Run job: 22 Run job: 24
以上內容是小編給大家介紹的Python使用multiprocessing實現一個最簡單的分布式作業調度系統,希望對大家有所幫助!
聲明:本網頁內容旨在傳播知識,若有侵權等問題請及時與本網聯系,我們將在第一時間刪除處理。TEL:177 7030 7066 E-MAIL:11247931@qq.com