Part 11 · 2 chapters · ~12 min
Build: Extension, FDW and Decoding Plugin
The capstone: a working extension with a custom type, operators and index support, a foreign data wrapper over a CSV file or HTTP API, and a logical decoding output plugin that emits changes as JSON lines.
29
A custom type with operators and B-tree support
code
-- currency_amount: (minor units, currency) with ordering only within a currency CREATE TYPE currency_amount AS (minor bigint, cur char(3)); CREATE FUNCTION ca_lt(a currency_amount, b currency_amount) RETURNS bool LANGUAGE sql IMMUTABLE AS $$ SELECT CASE WHEN a.cur <> b.cur THEN (a.cur < b.cur) ELSE a.minor < b.minor END $$; -- ca_le, ca_eq, ca_ge, ca_gt similarly; ca_cmp returns -1, 0, 1 CREATE OPERATOR < (LEFTARG = currency_amount, RIGHTARG = currency_amount, FUNCTION = ca_lt, COMMUTATOR = >, NEGATOR = >=); CREATE OPERATOR CLASS currency_amount_ops DEFAULT FOR TYPE currency_amount USING btree AS OPERATOR 1 <, OPERATOR 2 <=, OPERATOR 3 =, OPERATOR 4 >=, OPERATOR 5 >, FUNCTION 1 ca_cmp(currency_amount, currency_amount); CREATE INDEX ON payouts (amount); -- now a B-tree index on your type
Package it as an extension (control file, versioned SQL script, PGXS Makefile from part 7), move the comparison functions to C for speed, and add a regression test suite with pg_regress (REGRESS = money in the Makefile).
30
An FDW and a decoding plugin
code
// the FDW callbacks you implement (C), simplified
GetForeignRelSize // estimate rows
GetForeignPaths // offer a ForeignPath with a cost
GetForeignPlan // build the ForeignScan node (pushdown decisions)
BeginForeignScan // open the CSV file / start the HTTP request
IterateForeignScan // return the next row as a tuple, or empty when done
EndForeignScan // close
-- use it
CREATE SERVER rates_api FOREIGN DATA WRAPPER http_fdw OPTIONS (url 'https://rates.example.com');
CREATE FOREIGN TABLE live_rates (pair text, rate numeric) SERVER rates_api;
SELECT * FROM live_rates WHERE pair = 'NGN/USD'; -- pushdown: ?pair=NGN/USD
// logical decoding output plugin callbacks
_PG_output_plugin_init(cb) { cb->begin_cb = ...; cb->change_cb = ...; cb->commit_cb = ...; }
// change_cb receives (relation, ReorderBufferChange) → write one JSON line per row change
SELECT * FROM pg_create_logical_replication_slot('jsonl', 'jsonl_decoder');
SELECT data FROM pg_logical_slot_get_changes('jsonl', NULL, NULL);