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);