PostgreSQL 協定除了 COPY 的串流能力,還實作了非同步訊息與通知:連線一旦建立,即使客戶端閒置,伺服器也能主動推送訊息給它。
PostgreSQL 通知#
從伺服器流向客戶端的訊息應由客戶端處理——可能是伺服器要重啟,也可能是應用層訊息。基本用法:
yesql# listen channel;
LISTEN
yesql# notify channel, 'foo';
NOTIFY
Asynchronous notification "channel" with payload "foo"
received from server process with PID 40430.- 訊息可以來自另一條連線——開幾個 psql 實例試試看。
- Payload 可以是任意文字,上限 8kB,因此能傳遞豐富的訊息,例如 JSON 編碼的值。
PostgreSQL 事件發佈系統#
觸發器一章看到:用 trigger 維護計數快取雖然實現了事件驅動,卻扼殺了並行與擴展性。改良方案:讓 trigger 通知外部客戶端——一個以 listen 註冊訊息的常駐程式(daemon),每收到通知就視需要更新 twcache.counters。由於只有單一 daemon 在聽通知、更新快取,並行問題就此繞開。
先實作通知用的 trigger,用 psql 當測試客戶端:
begin;
create or replace function twcache.tg_notify_counters ()
returns trigger
language plpgsql
as $$
declare
channel text := TG_ARGV[0];
begin
PERFORM (
with payload(messageid, rts, favs) as
(
select NEW.messageid,
coalesce(
case NEW.action
when 'rt' then 1
when 'de-rt' then -1
end,
0
) as rts,
coalesce(
case NEW.action
when 'fav' then 1
when 'de-fav' then -1
end,
0
) as favs
)
select pg_notify(channel, row_to_json(payload)::text)
from payload
);
RETURN NULL;
end;
$$;
CREATE TRIGGER notify_counters
AFTER INSERT
ON tweet.activity
FOR EACH ROW
EXECUTE PROCEDURE twcache.tg_notify_counters('tweet.activity');
commit;測試:先 listen "tweet.activity";,再於單一語句中插入 7 筆活動(含重複的 (33, 'rt') 等)。結果只收到 4 則通知:
INSERT 0 7
Asynchronous notification "tweet.activity" with payload
"{"messageid":33,"rts":1,"favs":0}" received from
server process with PID 73216.
Asynchronous notification "tweet.activity" with payload
"{"messageid":33,"rts":-1,"favs":0}" received from ...
(共 4 則)7 筆 insert 只產生 4 則通知——這是
NOTIFY文件明載的行為:同一交易內對同一 channel 發出相同 payload 的多次通知,伺服器可以只送一則;payload 不同的通知則一定分開送達,不同交易的通知也絕不合併。除去重複折疊之外,NOTIFY 保證同一交易內的通知依發送順序送達,不同交易的通知依提交順序送達。
因此,我們以 notify 實作快取伺服器要正確,前提是主應用在單一交易內只發出彼此相異的 tweet.activity 動作。對我們的用途這不是致命傷。把測試改成一筆一交易後,四筆 insert 就得到預期的四則通知。
通知與快取維護#
有了伺服器端基礎設施(PostgreSQL 是伺服器、後端應用是客戶端),就能以事件驅動方式維護 twcache.counters。LISTEN/NOTIFY 非常適合維護快取——但因為通知只送給 notify 當下正在監聽的連線,快取維護服務必須嚴格依此順序運作:
- 連上目標資料庫並下
listen指令; - 從單一事實來源(single source of truth)取回目前值,重設快取;
- 處理陸續到來的通知、更新記憶體內快取,並依快取失效政策不時把記憶體快取同步(materialize)到落地位置。
快取本體可以放在維護服務內(例如同時處理通知並經 HTTP API 供應快取),也可以用 Memcached、Redis 等現成方案。本例用 Go 實作維護服務,快取本身放在 PostgreSQL 表:
create table twcache.counters
(
messageid bigint not null primary key,
rts bigint,
favs bigint
);服務執行時的日誌:
2017/09/21 22:00:36 Connecting to postgres:///yesql?sslmode=disable…
2017/09/21 22:00:36 Listening to notifications on channel "tweet.activity"
2017/09/21 22:00:37 Cache initialized with 6 entries.
2017/09/21 22:00:37 Start processing notifications, waiting for events…
2017/09/21 22:00:42 Received event: {"messageid":33,"rts":1,"favs":0}
...
2017/09/21 22:00:47 Materializing 6 events from memoryGo 客戶端全文 212 行放不進書頁,但 materialize 函式值得一看——它把記憶體內的快取結構整包推回 PostgreSQL:
type Counter struct {
MessageId int `json:"messageid"`
Rts int `json:"rts"`
Favs int `json:"favs"`
}
type Cache map[int]*Counter利用 Go 內建的 JSON 序列化,把整個快取變成單一 JSON 物件傳給 PostgreSQL(db.Query(q, js)),再由這條 SQL 處理:
with rec as
(
select rec.*
from json_each($1) as t,
json_populate_record(null::twcache.counters, value) as rec
)
insert into twcache.counters(messageid, rts, favs)
select messageid, rts, favs
from rec
on conflict (messageid)
do update
set rts = counters.rts + excluded.rts,
favs = counters.favs + excluded.favs
where counters.messageid = excluded.messageid- 撇開 JSON 技巧,序列化複雜資料結構到多列的傳統做法,見下一章的批次更新範例。
- 啟動時抓取快取初始值的查詢,直接沿用先前的 view:
select messageid, rts, favs from tweet.message_with_counters;。
延伸:trigger 成本的基準測試
trigger 處理當然有成本:
CL-USER> (concurrency::concurrency-test 100 100 35)
Starting benchmark for updates
Updating took 8.428939 seconds, did 10000 rts
Starting benchmark for inserts
Inserting took 10.351908 seconds, did 10000 rts注意這組數字已不能直接比較:trigger 裝在 tweet.activity 的 after insert 上,update 基準完全沒呼叫 trigger,insert 基準卻呼叫了一萬次 trigger 函式。至於並行方面,通知在提交時序列化——與 PostgreSQL commit log 的序列化方式相同——伺服器並無額外負擔。
維護服務收了 10,000 則 JSON 通知,只對快取表回報一次累計值。查看快取表(table twcache.counters;)可見 messageid 1–6 的既有數值也在——那是服務啟動時「先從資料庫真值重設快取」規則的結果,畢竟先前那幾輪測試時快取服務根本還沒寫出來。
Listen 與 Notify 的限制#
使用 PostgreSQL 通知功能的應用必須能承受漏接事件:通知只送給當下連線中的客戶端。任何佇列(queueing)機制都要求「沒有 worker 連線時累積的事件」保留到下次連線——複寫(replication)正是事件佇列的特例——而 listen/notify 無法正確實作佇列。快取維護服務恰是這個功能的完美用例,因為服務啟動或重啟時重設快取很容易。
驅動程式的支援#
listen/notify 的支援程度依驅動而異:
- Java JDBC:關鍵限制是無法接收非同步通知,必須輪詢(poll)後端;poll 可設 timeout,但其他執行緒的語句執行會被擋住。官方文件(PostgreSQL Extensions to the JDBC API)附有完整的類別實作範例。
- Python(Psycopg):支援避免忙等(busy looping)的進階技巧——與其定時輪詢,不如用
select()這類 I/O completion 函式睡到核心喚醒,沒資料可讀時完全不耗 CPU。 - Go(pq):支援通知且不需輪詢——本章範例用的正是它。
- 其他語言請查閱所用驅動的文件。