пятница, 18 марта 2016 г.

30. ZeroMQ и Delphi. ZMQ 4.2: threas-safe сокеты ZMQ_CLIENT и ZMQ_SERVER. Новая модель связи Client-Server.

ZMQ 4.2: новые thread-safe сокеты ZMQ_CLIENT и ZMQ_SERVER, новая модель Client-Server



  Модель Client-Server обеспечивает общение одного сокета ZMQ_SERVER с одним (или множество) сокетами ZMQ_CLIENT. Клиент начинает общение, далее общение между клиентом и сервером происходит в асинхронном режиме, как между равнозначными партнерами.

Новая модель предназначена для замены старой модели Rлиент-Cервер с использованием пар сокетов ZMQ_DEALER/ZMQ_ROUTER и пар ZMQ_REQ/ZMQ_REP для схемы Запрос-Ответ.
Формальное описание нового шаблона можно найти здесь.

понедельник, 7 декабря 2015 г.

23.3 ZeroMQ: Пересылка файлов. Модель №3 - клиент использует конвейерное управление потоком на основе кредитования сервера запросами.

Начало - здесь.

Сервер может отсылать по 10 кусков за раз, затем ждать однократного подтверждения.
Это, в общем, бессмысленно: все равно что умножать размер куска на 10.
Сервер может отправлять клиенту порции файла без перезапросов от клиента, периодически делая небольшие паузы, чтобы загрузить сеть настолько, насколько она может справиться. Для этого сервер должен знать, что происходит с сетью. Получение такой информации представляется непростой задачей. Кроме того, непонятно, что делать, когда сеть быстрая, а клиенты медленные. Кто будет заниматься буферизацией сообщений?
Сервер мог бы отслеживать состояние исходящей очереди и отправлять сообщения тогда, когда в очереди есть свободное место. Но ZeroMQ не позволяет таких вольностей. И сервер, и сеть могут оказаться достаточно быстрыми, а клиент - оказаться маленьким медленным устройством.
Можно, в конце концов, модифицировать libzmq так, чтобы изменить поведение при достижении границы HWM. Возможно, нужно блокировать новые сообщения. Тоже не очень здорово: один маленький медленный клиент заблокирует весь сервер.
Можно попробовать возвращать клиенту сообщение с ошибкой. Тогда усложнится сервер. Ох. Пока самое лучшее, что может сделать сервер - это отбросить сообщение.
В общем, все эти варианты либо усложняют логику, либо вообще легко приводят систему в нерабочее состояние
Все, что нам нужно - это дать возможность клиенту возможность сообщать серверу о своей готовности к работе. Если мы все сделаем правильно, данные будут поступать к клиенту непрерывным потоком, но лишь тогда, когда клиент будет готов их принять.



23.2 ZeroMQ: Пересылка файлов. Модель №2 - клиент запрашивает каждую порцию файла по отдельности.


Продолжение (начало - здесь).

Вторая версия протокола пересылки файла предусматривает, что клиент каждый раз запрашивает по одной порции файла, а сервер, соответственно, возвращает по одной порции на каждый запрос, полученный от клиента:



пятница, 4 декабря 2015 г.

23.1 ZeroMQ: Пересылка файлов. Модель №1 - отправляем сразу весь файл большими порциями.


Задача: переслать файл.


ZMQ "искаропки" отлично справляется с пересылкой сообщений и оповещениями, но вот с пересылкой файлов все не так очевидно. Можно, конечно, тупо попробовать засунуть файл в сообщение целиком.
Однако, есть возражения:
1. Оперативная память, к сожалению, не резиновая.
2. Ладно, отправили по сети файл-сообщение. А сеть оказалась медленной (какой-нибудь древний вай-фай), да еще и неустойчивой. ZMQ, как мы помним, пытается обеспечить доставку сообщения даже в случае сбоев. Итак, несколько попыток пересылки по 1 гигабайту - здорово?
3. А если этот гигантский файл хотят скачать несколько клиентов одновременно? Каждому клиенту выделить буфер нужного размера, отправить и ждать? Свободная память закончится гораздо быстрее. То есть, "правильный" протокол передачи файлов должен учитывать ограниченность размеров ОЗУ.
...
...начинаем строить правильный протокол передачи файлов.

суббота, 10 января 2015 г.

22.6.1 ZeroMQ: надежные схемы "Запрос/Ответ". Сервис - ориентированная надежная очередь. Шаблон "Мажордом". "Объектный" стиль кода.

(Начало).

Ничего нового в нижеследующем коде нет, за исключением того, что он написан с использованием объектов и классов Delphi.

Код клиента в "объектном" стиле::

  
program mdclient;

{$APPTYPE CONSOLE}

uses
  SysUtils,
  zmq_h,
  czmq_h,
  mdcliapi in 'mdcliapi.pas',
  mdp in 'mdp.pas';

procedure DoIt;
var
  i: Integer;
  fSession: TmdClient;
  fVerbose: Boolean;
  reply: p_zmsg_t;
  request: p_zmsg_t;
begin
  fVerbose := (ParamCount > 0) and (ParamStr(1) = '-v');
  Writeln(ParamStr(1));
  fSession := TmdClient.Create('tcp://localhost:5555', fVerbose);
  try

    i := 0;
    for i := 0 to 99999 do begin
      request := zmsg_new();
      zmsg_pushstr(request, 'Hello world');
      reply := fSession.Send('echo', request);
      if reply <> nil then
        zmsg_destroy(reply)
      else
        break; //  Прерывание или отказ
    end;

    zclock_log('%d requests / replies processed'#10, i);
  finally
    fSession.Free;
  end;

end;
begin
  DoIt;
  Readln;
end.


Код API клиента в "объектном" стиле::

  

unit mdcliapi;

interface
uses
  mdp
  , zmq_h
  , czmq_h
  , ZMQ_Utils;


type

  TmdClient = class(TObject)
  private
    Ctx: p_zctx_t; //  ZMQ контекст
    Broker: string;
    sctClient: Pointer; //  Сокет для связи с брокером
    Verbose: Boolean; //  Протоколирование в stdout
    Timeout: integer; //  Таймаут запроса
    Retries: integer; //  Число попыток

    procedure ConnectToBrocker;

  public
    constructor Create(aBroker: string; aVerbose: Boolean);
    destructor Destroy; override;
    function Send(service: PChar; var request: p_zmsg_t): p_zmsg_t;

  end;


implementation

uses
  SysUtils;


procedure TmdClient.ConnectToBrocker;
begin
  if sctClient <> nil then
    zsocket_destroy(ctx, sctClient);
  sctClient := zsocket_new(ctx, ZMQ_REQ);
  zmq_connect(sctClient, PChar(broker));
  if verbose then
    zclock_log('I: connecting to broker at %s...', Broker);
end;

constructor TmdClient.Create(aBroker: string; aVerbose: Boolean);
begin
  Ctx := zctx_new();
  Broker := aBroker;
  Verbose := aVerbose;
  Timeout := 2500; //  msecs
  Retries := 3; //  Before we abandon
  ConnectToBrocker();
end;

destructor TmdClient.Destroy;
begin
  inherited;
  zctx_destroy(Ctx);
  Broker := '';
end;



function TmdClient.Send(service: PChar; var request: p_zmsg_t): p_zmsg_t;
var
  header: p_zframe_t;
  item: zmq_pollitem_t;
  msg: p_zmsg_t;
  rc: Integer;
  reply_service: p_zframe_t;
  retries_left: Integer;
begin
  assert(request <> nil);

  //  В соответствии с требованиями протокола, должны быть фреймы-префиксы
  //  Фрейм 1: "MDPCxy" (шесть байт, MDP/Client x.y)
  //  Фрейм 2: Имя сервиса (печатная строка)
  zmsg_pushstr(request, service);
  zmsg_pushstr(request, cMDPC_CLIENT);
  if verbose then begin
    zclock_log('I: send request to ''%s'' service:', service);
    zmsg_print(request);
  end;
  retries_left := Retries;
  while (retries_left > 0) and (zctx_interrupted = 0) do begin
    msg := zmsg_dup(request); // Работаем с копией (вдруг оригинал пригодится?)
    zmsg_send(msg, sctClient);

    zPollItemInit(item, sctClient, 0, ZMQ_POLLIN, 0);

    //  При любом блокирующем вызове (libzmq), в случае ошибки  будет -1;
    //  Теоретически, следует учесть разные коды завершения,
    //  но на практике достаточно учесть {EINTR} (Ctrl-C):

    rc := zmq_poll(@item, 1, Timeout * ZMQ_POLL_MSEC);
    if rc = -1 then
      break; //  Прерван

    //  Что-то приняли, обрабатываем
    if (item.revents and ZMQ_POLLIN) <> 0 then begin
      msg := zmsg_recv(sctClient);
      if Verbose then begin
        zclock_log('I: received reply:');
        zmsg_print(msg);
      end;
      //  В реальном коде обрабботка отказа должна быть более тщательной
      assert(zmsg_size(msg) >= 3);

      header := zmsg_pop(msg);
      assert(zframe_streq(header, cMDPC_CLIENT));
      zframe_destroy(header);

      reply_service := zmsg_pop(msg);
      assert(zframe_streq(reply_service, service));
      zframe_destroy(&reply_service);

      zmsg_destroy(request);
      Result := msg; //  Успех
      Exit;
    end
    else if s_Dec(retries_left) > 0 then begin
      if (Verbose) then
        zclock_log('W: no reply, reconnecting...');
      ConnectToBrocker;
    end
    else begin
      if (Verbose) then
        zclock_log('W: permanent error, abandoning');
      break; //  Всё

    end;
  end;
  if zctx_interrupted <> 0 then
    zclock_log('W: interrupt received, killing client...'#10);
  zmsg_destroy(request);
  Result := nil;
end;


end.

Код Рабочего в "объектном" стиле::

  
program mdworker;

{$APPTYPE CONSOLE}

uses
  SysUtils,
  zmq_h,
  czmq_h,
  mdwrkapi in 'mdwrkapi.pas',
  mdp in 'mdp.pas';

procedure DoIt;
var
  fVerbose: Boolean;
  fReply: p_zmsg_t;
  fRequest: p_zmsg_t;
  fSession: TmdWorker;
begin
  fVerbose := (ParamCount > 0) and (ParamStr(1) = '-v');
  fSession := TmdWorker.Create('tcp://localhost:5555', 'echo', fVerbose);
  fReply := nil;
  while (true) do begin
    fRequest := fSession.Recv(fReply);
    if (fRequest = nil) then
      break; //  Рабочий был прерван
    fReply := fRequest; //  Эхо ... :-)
  end;
  fSession.Free;
  Exit;

end;
begin
  doIt;
  Readln;
end.


Код API рабочего в "объектном" стиле::

  

unit mdwrkapi;

interface
uses
  zmq_h
  , czmq_h
  ;
type
  TmdWorker = class(TObject)
    Ctx: p_zctx_t; //  Наш контекст
    Broker: string;
    Service: string;
    sctWorker: Pointer; //  Сокет для связи с брокером Socket to broker
    Verbose: Boolean; //  Протоколировать действия в stdout

    //  Управление хартбитингом
    HeartbeatAt: uint64_t;
    Liveness: size_t; //  Сколько осталось попыток
    Heartbeat: integer; //  Задержка хартбитинга, в миллисекундах
    Reconnect: Integer; //  Задержка реконнекта в миллисекундах

    ExpectReply: Boolean;
    ReplyTo: p_zframe_t;
  private
    procedure ConnectToBroker;

    procedure SendToBroker(aCommand, aOption: PChar; aMsg: p_zmsg_t);
  public
    constructor Create(const aBroker, aService: string; aVerbose: Boolean);
    destructor Destroy; override;
    function Recv(var reply: p_zmsg_t): p_zmsg_t;


  end;
const
  //  Параметры надежности
  cHEARTBEAT_LIVENESS = 3; //  3-5 достаточно


implementation

uses
  ZMQ_Utils, mdp;


constructor TmdWorker.Create(const aBroker, aService: string; aVerbose: Boolean);
begin
  assert(aBroker <> '');
  assert(aService <> '');

  Ctx := zctx_new();
  Broker := aBroker;
  Service := aService;
  Verbose := aVerbose;
  Heartbeat := 2500; //  msecs
  Reconnect := 2500; //  msecs
  ConnectToBroker();

end;

destructor TmdWorker.Destroy;
begin
  inherited;
  zctx_destroy(Ctx);
  Broker := '';
  Service := '';
end;

function TmdWorker.Recv(var reply: p_zmsg_t): p_zmsg_t;
//  Метод сперва отсылает брокеру ответ, а затем ждет нового запроса.
//  Название метода странное, угу.

//  Отправляет ответ брокеру, если он есть и ждет следующего запроса


var
  fCommand: p_zframe_t;
  fEmpty: p_zframe_t;
  fHeader: p_zframe_t;
  fItem: zmq_pollitem_t;
  fMsg: p_zmsg_t;
  RC: Integer;
begin
  //  Формирование и отсылка ответа (reply), если он пустой
  assert((reply <> nil) or (not ExpectReply));
  if (reply <> nil) then begin
    assert(ReplyTo <> nil);
    zmsg_wrap(reply, ReplyTo);
    SendToBroker(cMDPW_REPLY, nil, reply);
    zmsg_destroy(reply);
  end;

  ExpectReply := true;

  while (true) do begin


    zPollItemInit(fItem, sctWorker, 0, ZMQ_POLLIN, 0);
    RC := zmq_poll(@fItem, 1, Heartbeat * ZMQ_POLL_MSEC);
    if RC = -1 then
      break; //  Прерван

    if (fItem.revents and ZMQ_POLLIN) <> 0 then begin
      fMsg := zmsg_recv(sctWorker);
      if fMsg = nil then
        break; //  Прерван
      if Verbose then begin
        zclock_log('I: received message from broker:');
        zmsg_print(fMsg);
      end;

      Liveness := cHEARTBEAT_LIVENESS;

      //  Ошибки не обрабатываем, просто шумим assert-от
      assert(zmsg_size(fMsg) >= 3);

      fEmpty := zmsg_pop(fMsg);
      assert(zframe_streq(fEmpty, ''));
      zframe_destroy(fEmpty);

      fHeader := zmsg_pop(fMsg);
      assert(zframe_streq(fHeader, cMDPW_WORKER));
      zframe_destroy(fHeader);

      fCommand := zmsg_pop(fMsg);
      if (zframe_streq(fCommand, cMDPW_REQUEST)) then begin
        //  Вообще-то, мы должны извлечь и сохранить все адреса вплоть до
        // пустого, но сейчас сохранем просто один...
        ReplyTo := zmsg_unwrap(fMsg);
        zframe_destroy(fCommand);

        //  В этом месте мы должны обрабатывать сообщение
        //  Сейчас просто возвращаем его в приложение

        Result := fMsg; //  Запрос на обработку
        Exit;

      end
      else
        if (zframe_streq(fCommand, cMDPW_HEARTBEAT)) then
          // ;               //  В случае хартбита ничего не делаем
        else
          if (zframe_streq(fCommand, cMDPW_DISCONNECT)) then
            ConnectToBroker()
          else begin
            zclock_log('E: invalid input message');
            zmsg_print(fMsg);
          end;
      zframe_destroy(&fCommand);
      zmsg_destroy(&fMsg);

    end
    else
      if s_Dec(Liveness) = 0 then begin
        if Verbose then
          zclock_log('W: disconnected from broker - retrying...');
        zclock_sleep(Reconnect);
        ConnectToBroker();
      end;
    //  Отправить хратбит, если пора
    if zclock_time() > HeartbeatAt then begin
      SendToBroker(cMDPW_HEARTBEAT, nil, nil);
      HeartbeatAt := zclock_time() + Heartbeat;
    end;

  end;
  if zctx_interrupted() <> 0 then
    z_log('W: interrupt received, killing worker...'#10);
  Result := nil;
end;

procedure TmdWorker.ConnectToBroker;
//  Коннект или реконнект к брокеру
begin
  if sctWorker <> nil then // Полный дисконнект, с разрушенеим сокета
    zsocket_destroy(Ctx, sctWorker);
  sctWorker := zsocket_new(Ctx, ZMQ_DEALER);
  zmq_connect(sctWorker, PChar(Broker));
  if Verbose then
    zclock_log('I: connecting to broker at %s...', PChar(Broker));

  //  Регистрация сервиса в брокере
  SendToBroker(cMDPW_READY, PChar(Service), nil);

  //  Если число попыток (liveness) обнулилось, считаем, что брокер отключен
  Liveness := cHEARTBEAT_LIVENESS;
  HeartbeatAt := zclock_time() + heartbeat;

end;


procedure TmdWorker.SendToBroker(aCommand, aOption: PChar; aMsg: p_zmsg_t);
//  Отправка сообщения брокеру
//  Если сообщения нет (пусто), создать его самостоятельно
begin
  if aMsg <> nil then
    aMsg := zmsg_dup(aMsg)
  else
    aMsg := zmsg_new();

  //  Формирование конверта сообщения в сответствии с протоколом
  if aOption <> nil then
    zmsg_pushstr(aMsg, aOption);
  zmsg_pushstr(aMsg, aCommand);
  zmsg_pushstr(aMsg, cMDPW_WORKER);
  zmsg_pushstr(aMsg, '');

  if Verbose then begin
    zclock_log('I: sending %s to broker',
      mdps_commands[Byte(aCommand^)]);
    zmsg_print(aMsg);

  end;
  zmsg_send(aMsg, sctWorker);
end;


end.


Если кто-то возжелает написать и брокер в объектном стиле, напоминаю, что вторым аргументом метода zhash_freefn() должна быть ссылка на функцию типа zhash_free_fn:
  
  zhash_free_fn = procedure(data: Pointer); cdecl;
Чтобы описать такой метод в рамках класса, следует использовать модификаторы class и static
  
    class procedure service_destroy(service: Pointer); static;
Честно говоря, я бы не стал заморачиваться с zhash_lookup(), а использовал бы вместо них, например, дельфийские TStringList с включенной сортировкой.

пятница, 9 января 2015 г.

22.6 ZeroMQ: надежные схемы "Запрос/Ответ". Сервис - ориентированная надежная очередь. Шаблон "Мажордом".

(Начало - здесь)

ZeroMQ: надежные схемы "Запрос/Ответ". Сервис - ориентированная надежная очередь. Шаблон "Мажордом".


Шаблон "Мажордом". Топология.


Прогресс ускоряется, когда в нем не участвуют юристы и разные комитеты. Одно-страничная спецификация MDP сразу  превращает PPP в нечто более солидное. Именно так и должна происходить разработка сложных архитектур: начинать следует с написания соглашений, а только потом писать софт, их реализующий.
Протокол "Мажордом" (The Majordomo Protocol - MDP) расширяет и углубляет PPP в одном интересном направлении: он добавляет "имя сервиса" к запросу, который отправляет клиент, и требует от рабочих регистрироваться для оказания конкретных услуг. Добавление имен сервиса превращает брокер из шаблона "Пират-параноик" в сервис - ориентированный брокер.

четверг, 8 января 2015 г.

22.5 ZeroMQ: надежные схемы "Запрос/Ответ". Хартбитинг в деталях.

  Запрос/Ответ. Надежные схемы "Запрос/Ответ". Хартбитинг в деталях.

Хартбитинг

 


Хартбитинг решает проблему распознавания "жив партнер или нет?" Эта проблема касается не только ZeroMQ. Протокол TCP имеет долгий таймаут (30 минут и больше), что делает невозможным определить - умер ли ли партнер, был ли дисконнект, или партнер просто уехал на выходные в Прагу пить водку.
Реализовать хартбитинг не так просто. Автор оригинала примеров для шаблона "Пират-параноик" пишет, что это заняло около пяти часов. Для реализации остальной части цепочки "запрос-ответ" потребовалось от силы десять минут. Особенно легко реализуются "ложные отказы", когда, к примеру, партнеры решают, что приключился дисконнект из-за неправильного хартбитинга.

В плане использования хартбитинга совместно с ZeroMQ разработчики обычно придерживаются трех подходов:

1. "А, и так сойдет!"


вторник, 6 января 2015 г.

22.4 ZeroMQ: надежные схемы "Запрос/Ответ". Высоконадежный брокер с очередью и хартбитингом (шаблон "Пират-параноик")


Запрос/Ответ. Высоконадежный брокер с очередью и хартбитингом (шаблон "Пират - параноик").

Шаблон "Пират - параноик"

 

 

Шаблон "Простой пират" достаточно хорош особенно тем, что он является простой комбинацией двух уже существующих шаблонов. Тем не менее, и у него тоже есть некоторые недостатки:

  • Он перестает работать в случае, когда очередь (брокер) падает и перезапускается. Клиент восстановится, а рабочие - нет. Хотя ZeroMQ и выполнит автоматический реконнект после перезапуска очереди, рабочие не пошлют сигнала "ГОТОВ" и, следовательно, не будут считаться доступными. Для исправления реализуем хартбитинг от очереди к рабочим так, чтобы рабочий смог определить, что очередь стала недоступной.
  • Очередь не обнаруживает отказов рабочих, поэтому, если рабочие падают во время простоя, брокер не может удалить таких рабочих их очереди доступных рабочих до тех пор, пока брокер не не пошлет такому рабочему запрос. Клиент будет ждать и делать перезапросы в никуда. Это не очень большая проблема, но это неприятно. Чтобы все работало правильно, реализуем хартбитинг от рабочего к очереди так, чтобы очередь могла определять потерянных рабочих на любом этапе.
Решим эти проблемы проблемы с помощью шаблона "Пират - параноик".

суббота, 27 декабря 2014 г.

22.3 ZeroMQ: надежные схемы "Запрос/Ответ". Надежность на основе прокси с очередью (шаблон "Простой пират")

  Запрос/Ответ. Надежность  на базе прокси с очередью (брокер с балансировкой нагрузки). Шаблон "Простой пират".

 Шаблон "Простой  Пират"

Расширим предыдущий шаблон обеспечения надежности типа "Ленивый пират", использовав прокси с очередями. Новый шаблон позволит общаться прозрачно с несколькими серверами, которые далее будем более точно называть "рабочими".

Во всех "пиратских"  шаблонах рабочие не сохраняют свое состояние. Разработка системы обмена сообщениями подразумевает, что мы ничего не знаем о том, что приложениям могут понадобится какие-либо разделяемые ресурсы, хранящие состояние, вроде баз данных. Использование прокси с очередями подразумевает, что рабочие могут приходить и уходить, не зная ничего о клиентах. Если падает один из рабочих, другой заменяет его. Это - хорошая, простая топология с одним слабым местом: а именно брокером с очередью в центре, который может стать проблемой для системы управления и единой точкой отказа.

четверг, 25 декабря 2014 г.

22.2 ZeroMQ: надежные схемы "Запрос/Ответ". Надежность на стороне клиента (шаблон "Ленивый пират")

  Запрос/Ответ. Надежность на стороне клиента. Шаблон "Ленивый пират".

 Шаблон "Ленивый  Пират"

 

 

Топология сети:




Небольшие изменения на стороне клиента позволяют получить надежную схему "Запрос-Ответ".
Вместо того, чтобы выполнять блокирующее чтение (receive) из сокета, поступаем иначе:

22.1 ZeroMQ: надежные схемы "Запрос/Ответ". Общие положения.

(Начало - здесь)

 Пора задуматься о надежности.


Будет рассмотрены повторяемые шаблоны, позволяющие добиться надежности при использовании схемы "Запрос-Ответ". Список шаблонов:
  • "Ленивый пират": надежная схема "запрос-ответ" на стороне клиента.
  • "Простой пират": надежная схема "запрос-ответ" с использованием балансировки нагрузки.
  • "Пират-параноик": надежная схема "запрос-ответ" с использованием хартбитинга.
  • "Мажордом": сервис - ориентированная надежная очередь.
  • "Титаник": диск-ориентированная надежная оффлайн очередь.
  • "Двоичная звезда": надежный отказоустойчивый бэкап - сервер.
  • "Фрилансер": надежная схема "запрос-ответ" без брокеров.

Что же такое - "Надежность"?

Ни о какой "теории надежности" речи не пойдет.

среда, 3 декабря 2014 г.

21.4 ZeroMQ: рабочий пример. Межброкерная маршрутизация. Вот что получилось.

(Начало - здесь).




Соберем все вместе.


Как и раньше, отдельный кластер будет представлен одним процессом.

Кода получилось довольно много, так как здесь объединяется логика обоих предыдущих примеров. Он достаточно хорошо имитирует кластеры, включающие клиентов и рабочих. 

Код:

21.3 ZeroMQ: рабочий пример. Межброкерная маршрутизация. Прототипирование локального и облачного потоков данных.


(Начало - здесь).

Прототипирование локального и облачного потоков.

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

Потоки данных задач:


Перед написанием кода (который все больше усложняется) обрисуем основы логики маршрутизации и разобьем её на простые, но надежные элементы.

21.2 ZeroMQ: рабочий пример. Межброкерная маршрутизация. Прототипирование потока данных о состоянии.


(Начало - здесь)

Прототипирование потока данных о состоянии


Так как в каждом потоке данных сокетов есть свои маленькие ловушки для раззяв, будем добавлять новые потоки данных в реальный код постепенно, шаг за шагом. Начнем с потока состояния.

Код прототипа обработчика потока состояния:

пятница, 28 ноября 2014 г.

21.1 ZeroMQ: рабочий пример. Межброкерная маршрутизация. Разбор задачи.

(Начало - здесь).


Постановка задачи.

Только что нам позвонил лучший клиент и попросил срочно спроектировать мощную систему облачных вычислений. В качестве облака он представляет себе совокупность множества датацентров, каждый кластер клиентов и рабочих которого работают вместе, как единое целое.

Нам эту задачу решить - как два байта переслать. 

Разберем задачу, используя ZeroMQ.

суббота, 15 ноября 2014 г.

20. ZeroMQ и Delphi: CZMQ, varargs и типы аргументов (zsock_bind(), zsock_connect(), zsock_send(), zsock_recv() и т.д.)


Использование функций с переменным числом аргументов.



В CZMQ появилось много методов с переменным число аргументов. В описании таких методов присутствует модификатор varargs: например:


function zsock_connect(p_self: p_zsock_t; format: PChar): Integer; varargs; 
  cdecl; external cZMQ_DllName;

или

function zsock_send(self: Pointer; picture: PChar): Integer; varargs; 
  cdecl; external cZMQ_DllName;


Как ими пользоваться?

вторник, 4 ноября 2014 г.

19. ZeroMQ: Асинхронный клиент-сервер.

(Начало - здесь).

Шаблон "Асинхронный клиент/сервер".

Будем создавать архитектуру сети N-1, когда несколько разных клиентов асинхронно общаются с одним сервером.

Работать это будет вот так:

- клиенты коннектятся к серверу и отправляют запросы;
- на каждый запрос сервер отправляет 0 или больше ответов;
- клиенты могут отправлять множество запросов без ожидания ответов;
- серверы могут отправлять множество ответов без ожидания новых запросов.

Топология:




18. ZeroMQ: реакторы CZMQ. ZLOOP.



(Начало - здесь).

Шаблон "Реактор".

<Теоория.>
Шаблон проектирования Реактор является событийно-управляемым шаблоном проектирования. Он предназначен для синхронной передачи запросов к сервису, поступающих параллельно от  одного или нескольких источников. Обработчик сервиса демультиплексирует  (разбирает) входящие запросы и  синхронно отправляет их ассоциированным обработчикам запросов.</Теория>

 Реактор напоминает процедуру окна Windows: разбирает поступившие сообщения и вызывает коллбеки, связанные с сообщениями.


Класс zloop - событийно-управляемый реактор.

понедельник, 3 ноября 2014 г.

17. ZeroMQ: брокер с балансировкой нагрузки + CZMQ. Обработка Ctrl+C в CZMQ. Реакторы.

(Начало - здесь).

Дальше рассмотрим кодирование брокера с балансировкой нагрузки с использованием API высокого уровня - CZMQ.

16. ZeroMQ: API высокого уровня. Библиотека CZMQ.

(Начало - здесь).

Наши примеры становятся все сложнее, код становится все более громозким и все менее располагает к пониманию.Вспомним, к примеру, последний вариант боркера с распределенной нагрузкой.

Громоздко, да. И это мы еще использовали наши вспомогательные процедуры вроде s_recv()/s_send(), а без них бы пришлось пересылать сообщения ZMQ и заниматься упаковкой-распоковкой данных.

В общем, если использовать низкуровневое API в сложных задачах, то и кодировать долго, и потом разбираться непросто.

Нужно переходить к более высоким уровням абстракции.
Например, для Delphi есть замечательная объектная библиотека: https://github.com/bvarga/delphizmq

Модуль zmqapi.pas предоставляет объектный интерфейс высокого уровня в соответствии видением прекрасного с создателя библиотеки и, надо полагать, в соответствии с задачами, которые стояли перед ним в момент написания.

К сожалению, больше года библиотека почти не обновляется, и зависла на поддержке ZeroMQ версий 2.* и 3.*.
~~~~~~~~~~~~

К счастью, выход есть: iMatrix (контора, которая и разрабатывает ZeroMQ) создала и развивает библиотеку API высокого уровня: http://czmq.zeromq.org/