Как распределить Uploads файлов через python?

Кто сможет подсказать как выполнить задачу ?

У меня есть файл, в нем находиться 100 строк, также у меня есть 3 сервера.

Из этих 100 строк мы отправляем по 10 строк на каждый сервер, и сохраняем их в файл, где каждый сервер обрабатывает свои строки. После того как он обработал первые 10x3 строк, у нас осталось 70 строк, и повторяем все по кругу , пока в файле из 100 строк не закончатся строки которые надо обработать.

Может кто нибудь показать пример самого распределения задачи ?


Ответы (3 шт):

Автор решения: DGDays

Ну, как вариант - сделать 3 потока. Далее каждый поток поочерёдно должен прочитать из файла 10 строк и записать их в отдельный временный файл, а из основного удалить. Ну и каждый поток выполняет эту процедуру до тех пор, пока строки не закончатся.

К примеру, можно открыть файл через open и далее прочитать с помощью .readlines(). Это создаст список всех строк и если мы оттуда удалим строки, то основной файл не пострадает. Удалять можно через del тут подробней. Ну а далее весь список из 10 значений отправлять на сервак. И так каждый поток, но только на свой сервер. Ну а дальше, думаю, не сложно написать программу по это алгоритму. Также желательно сделать очередь на чтение из списка, иначе два потока могут прочитать одни и теже строчки и одновременно их удалять, из-за чего у одного потока будет ошибка.

Ну, это примерный план по работе программы. Далее вы можете его сами модифицировать.

Ссылки:

Потоки

del для списка

Работа с файлами Python

→ Ссылка
Автор решения: Sergey

показать пример самого распределения задачи

Что такое "распределение задачи" я не знаю. Но, по сути, делал бы так:

  1. Разработал бы некий протокол обмена сообщениями между клиентом и сервером. Этот протокол содержит всего два типа сообщений: ТУДА (посылаем 10 строк на сервер) и ОБРАТНО (Сервер подтверждает, что он закончил обработку)
  2. Разработал бы программу сервера. Она очень простая: ждать прихода сообщения, обработать строки и послать сообщение об окончании обработки. В принципе, в протокол можно добавить некие дополнительные сообщения, например, команду прекращения работы сервера.
  3. Разработал бы программу клиента. Она чуть посложнее.

На псевдокоде:

  1. Отправить 10 строк серверу № 1
  2. Отправить 10 строк серверу № 2
  3. Отправить 10 строк серверу № 3
  4. цикл пока (Есть не отправленные строки)
  5. Ждём прихода сообщения от сервера
    
  6. Определяем - от кого пришло подтверждение
    
  7. Этому серверу отправляем очередные 10 строк.
    

Я надеюсь, как работать с сокетами в python - объяснять не нужно?

→ Ссылка
Автор решения: eri

Чтоб утилизировать сервера по полной надо отправлять им задачи по мере их выполнения. Вот пример через очередь и потоки.

Несколько импортов:

from multiprocessing.dummy import Pool 
import functools

тут обычные треды, но мне апи multiprocessing нравится больше. Можно всё это сделать через массив с threading.Thread, но зачем, если все уже красиво сделанно в multiprocessing.dummy. Плюс легкий переход процессы, если задача станет затратной по процессору.

Входные данные и функция для загрузки и получения результат

servers = ['http://server1.org', 'http://server2.org', 'http://server3.org']
f = open('file.txt', 'r')


def upload(server, chunk):
    pass # реализуйте аплоад тут и дождитесь завершения обработки

А вот отправка, пока в очереди есть задания - отправляем их. В конце пришлем индикатор конца очереди

def sender(q, server):
    while True:
        chunk = q.get()
        if not chunk:
            break
        upload(server, chunk)

Функция для нарезки файла по 10 строк

def chunker(q, f):
   chunk = []
   for line in f:
       chunk.append(line)
       if len(chunk) == 10:
          q.put(''.join(chunk))
          chunk = []
   q.put(''.join(chunk))

И запускаем. pool.map_async выполняет задачу в отдельном потоке, а нарезка в основном.

q = queue.Queue(10) # очередь на 10 по 10 строк
pool = Pool()
pool.map_async(functools.partial(sender, q), servers)
chunker(f, q)

Индикатор конца очереди и дождаться завержения работы серверов:

for _ in range(len(servers)):
    q.put(None)

pool.join()
→ Ссылка