コンテンツにスキップ

プロセスグループ

プロセスを名前付きグループに参加させ、クラスタ全体のすべてのメンバーにブロードキャストします。Erlang/OTP pg をモデルにしています: グループは動的で、プロセスは複数のグループに所属でき、メンバーシップはクラスタ全体で追跡され、最終的整合性があります。

スコープエントリ種別とその設定についてはプロセスグループを参照。クラスタリングモデル全体についてはクラスタガイドを参照。

local pg = require("pg")

プロセスグループはスコープpg.scope レジストリエントリ — の中に存在します。インスタンスを取得するにはそれを開きます:

local group, err = pg.open("app:pg")
if err then
return nil, err
end
パラメータ説明
idstringスコープエントリID(形式: "namespace:name"

戻り値: pg.Instance, error

権限: スコープ id に対する pg.open

インスタンスはプロセス終了時に自動的に解放されます。早期に解放するには release() を呼び出します。他の操作はすべてインスタンスのメソッドで、: で呼び出します。

local ok, err = group:join("workers") -- 単一グループ
local ok, err = group:join({"workers", "all"}) -- バッチ
local ok, err = group:leave("workers")
パラメータ説明
groupstring | string[]グループ名、またはバッチ操作の名前リスト

戻り値: boolean, error

プロセスは同じグループに複数回参加できます。完全に離脱するには同じ回数 leave する必要があります(マルチ参加セマンティクス)。leave はバッチ全体でベストエフォートで、指定されたいずれのグループにもメンバーでない場合のみエラーを返します。

権限: 各グループ名に対する pg.join / pg.leave

local members, err = group:get_members("workers") -- 全ノード
local local_members, err = group:get_local_members("workers") -- このノードのみ
パラメータ説明
groupstringグループ名

戻り値: string[], error — PID文字列の配列(不明なグループは空)

権限: グループ名に対する pg.get_members / pg.get_local_members

local groups, err = group:which_groups() -- クラスタ内の全グループ
local local_groups, err = group:which_local_groups() -- ローカルメンバーを持つグループ

戻り値: string[], error — 現在少なくとも1つのメンバーを持つグループ名

権限: pg.which_groups / pg.which_local_groups

グループのすべてのメンバーにメッセージを送信します。各メンバーは呼び出しプロセスから topic 名でメッセージを受け取ります — process.listen(topic) で処理します。

local ok, err = group:broadcast("workers", "task", {id = 42}) -- 全ノード
local ok, err = group:broadcast_local("workers", "task", {id = 42}) -- このノードのみ
パラメータ説明
groupstring対象グループ
topicstringメッセージトピック
...anyゼロ個以上のペイロード値

戻り値: boolean, error

権限: グループ名に対する pg.broadcast / pg.broadcast_local

monitor は1つのグループの参加/離脱イベントをサブスクライブし、現在のメンバーをアトミックに返します — サブスクリプションとスナップショットの間でメンバーシップ変更が抜け落ちることはありません。

local sub, members, err = group:monitor("workers")
if err then
return nil, err
end
for _, pid in ipairs(members) do
-- サブスクリプション時の現在のメンバー
end
local ch = sub:channel()
local event = ch:receive() -- {kind = "member.joined" | "member.left", path = "workers", data = {...}}
sub:close() -- アンサブスクライブ。sub:close({flush = true}) でキューされたイベントを先にドレイン
パラメータ説明
groupstring監視するグループ

戻り値: pg.Subscription, string[], error — サブスクリプションと現在のメンバーのスナップショット

権限: グループ名に対する pg.monitor

events はスコープ内のすべてのグループにまたがるメンバーシップ変更をサブスクライブし、すべてのグループとそのメンバーのスナップショットを返します。

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

戻り値: pg.Subscription, table, error

権限: pg.events

サブスクリプションチャネルで配信されるイベントには以下が含まれます:

フィールド説明
systemstring常に "pg"
kindstring"member.joined" または "member.left"
pathstringグループ名
datatable{Group = string, PIDs = string[]} — 影響を受けるメンバー

サブスクリプションチャネルはバッファ付き(容量64)。遅いコンシューマがバッファを満たすと、そのサブスクリプションへのイベントはドロップされます。

group:release()

インスタンスを即座に解放します。冪等です。解放後はすべてのメソッドがエラーを返します。プロセス終了時にもクリーンアップは自動的に実行されます。

戻り値: boolean

権限メソッドリソース
pg.openpg.open()scope id
pg.joinjoin()group name
pg.leaveleave()group name
pg.get_membersget_members()group name
pg.get_local_membersget_local_members()group name
pg.which_groupswhich_groups()(scope)
pg.which_local_groupswhich_local_groups()(scope)
pg.broadcastbroadcast()group name
pg.broadcast_localbroadcast_local()group name
pg.monitormonitor()group name
pg.eventsevents()(scope)
条件種別
権限拒否errors.PERMISSION_DENIED
引数が欠損または空errors.INVALID
スコープが見つからないerrors.NOT_FOUND
メンバーでないグループからの離脱errors.INVALID
インスタンスが解放済みerrors.INVALID

エラーの処理についてはエラー処理を参照。