Pular para o conteúdo

Grupos de Processos

Agrupe processos em grupos nomeados e faça broadcast para todos os membros em todo o cluster. Modelado no pg do Erlang/OTP: grupos são dinâmicos, um processo pode pertencer a muitos grupos e a associação é rastreada em todo o cluster e é eventualmente consistente.

Para o tipo de entrada de escopo e sua configuração, veja Grupos de Processos. Para o modelo de clustering mais amplo, veja o Guia de Cluster.

local pg = require("pg")

Um grupo de processos vive dentro de um escopo — uma entrada de registro pg.scope. Abra-o para obter uma instância sobre a qual você opera:

local group, err = pg.open("app:pg")
if err then
return nil, err
end
ParâmetroTipoDescrição
idstringID da entrada de escopo (formato: "namespace:name")

Retorna: pg.Instance, error

Permissão: pg.open no id do escopo

A instância é liberada automaticamente quando o processo sai; chame release() para liberá-la antes. Todas as outras operações são métodos na instância, chamados com :.

local ok, err = group:join("workers") -- grupo único
local ok, err = group:join({"workers", "all"}) -- lote
local ok, err = group:leave("workers")
ParâmetroTipoDescrição
groupstring | string[]Nome do grupo, ou lista de nomes para operação em lote

Retorna: boolean, error

Um processo pode entrar no mesmo grupo mais de uma vez; deve sair o mesmo número de vezes para partir completamente (semântica multi-join). leave é best-effort em um lote e retorna erro apenas quando o processo não era membro de nenhum dos grupos nomeados.

Permissões: pg.join / pg.leave em cada nome de grupo

local members, err = group:get_members("workers") -- todos os nós
local local_members, err = group:get_local_members("workers") -- apenas este nó
ParâmetroTipoDescrição
groupstringNome do grupo

Retorna: string[], error — array de strings PID (vazio para grupo desconhecido)

Permissões: pg.get_members / pg.get_local_members no nome do grupo

local groups, err = group:which_groups() -- todos os grupos no cluster
local local_groups, err = group:which_local_groups() -- grupos com membro local

Retorna: string[], error — nomes de grupos que atualmente têm pelo menos um membro

Permissões: pg.which_groups / pg.which_local_groups

Envia uma mensagem para todos os membros de um grupo. Cada membro a recebe sob topic do processo chamador — trate com process.listen(topic).

local ok, err = group:broadcast("workers", "task", {id = 42}) -- todos os nós
local ok, err = group:broadcast_local("workers", "task", {id = 42}) -- apenas este nó
ParâmetroTipoDescrição
groupstringGrupo alvo
topicstringTópico da mensagem
...anyZero ou mais valores de payload

Retorna: boolean, error

Permissões: pg.broadcast / pg.broadcast_local no nome do grupo

monitor inscreve-se em eventos de entrada/saída para um grupo e retorna os membros atuais atomicamente — nenhuma mudança de associação pode ocorrer entre o snapshot e a inscrição.

local sub, members, err = group:monitor("workers")
if err then
return nil, err
end
for _, pid in ipairs(members) do
-- membros atuais no momento da inscrição
end
local ch = sub:channel()
local event = ch:receive() -- {kind = "member.joined" | "member.left", path = "workers", data = {...}}
sub:close() -- cancelar inscrição; sub:close({flush = true}) drena eventos enfileirados primeiro
ParâmetroTipoDescrição
groupstringGrupo a observar

Retorna: pg.Subscription, string[], error — a inscrição e um snapshot dos membros atuais

Permissão: pg.monitor no nome do grupo

events inscreve-se em mudanças de associação em todos os grupos do escopo e retorna um snapshot de todos os grupos com seus membros.

local sub, snapshot, err = group:events()
-- snapshot: { ["workers"] = {pid, ...}, ["all"] = {pid, ...} }
local event = sub:channel():receive()
sub:close()

Retorna: pg.Subscription, table, error

Permissão: pg.events

Eventos entregues em um channel de inscrição contêm:

CampoTipoDescrição
systemstringSempre "pg"
kindstring"member.joined" ou "member.left"
pathstringO nome do grupo
datatable{Group = string, PIDs = string[]} — os membros afetados

Channels de inscrição têm buffer (capacidade 64); se um consumidor lento encher o buffer, eventos adicionais para essa inscrição são descartados.

group:release()

Libera a instância imediatamente. Idempotente; após a liberação, cada método retorna um erro. A limpeza também ocorre automaticamente quando o processo sai.

Retorna: boolean

PermissãoMétodoRecurso
pg.openpg.open()id do escopo
pg.joinjoin()nome do grupo
pg.leaveleave()nome do grupo
pg.get_membersget_members()nome do grupo
pg.get_local_membersget_local_members()nome do grupo
pg.which_groupswhich_groups()(escopo)
pg.which_local_groupswhich_local_groups()(escopo)
pg.broadcastbroadcast()nome do grupo
pg.broadcast_localbroadcast_local()nome do grupo
pg.monitormonitor()nome do grupo
pg.eventsevents()(escopo)
CondiçãoTipo
Permissão negadaerrors.PERMISSION_DENIED
Argumento ausente ou vazioerrors.INVALID
Escopo não encontradoerrors.NOT_FOUND
Sair de um grupo sem associaçãoerrors.INVALID
Instância liberadaerrors.INVALID

Veja Error Handling para trabalhar com erros.