Can pg_notify be as fast inside a trigger as outside?

Viewed 35

I noticed that statements like SELECT pg_notify('foo', 'bar') are incredibly quick to execute.

BenchmarkNotify-8                         257728             54542 ns/op

Then simple updates to random rows in a table are ~5x slower.

Init statement:

CREATE TABLE table1 (
    id SERIAL PRIMARY KEY,
    int INT
);

Benchmarked statement:

INSERT INTO table1 (id, int) VALUES($1, $2)
ON CONFLICT (id) DO UPDATE
SET int = $3
BenchmarkUpdate-8                          44913            289502 ns/op

But updates to a table with a trigger that runs PERFORM pg_notify('foo', 'bar') are ~5x slower still.

Init statement:

CREATE TABLE table1 (
    id SERIAL PRIMARY KEY,
    int INT
);

CREATE OR REPLACE FUNCTION table1_fn() RETURNS TRIGGER AS $$
BEGIN
    PERFORM pg_notify('foo', 'bar');
    RETURN NEW;
END;
$$ LANGUAGE plpgsql;

CREATE TRIGGER table1_trg AFTER UPDATE ON table1 FOR EACH ROW EXECUTE PROCEDURE table1_fn();

Benchmarked statement:

INSERT INTO table1 (id, int) VALUES($1, $2)
ON CONFLICT (id) DO UPDATE
SET int = $3
BenchmarkUpdateTriggerNotify-8              8110           1502632 ns/op

I'd like to understand why this slowdown is so large and whether it can be avoided? I would expect BenchmarkUpdateTriggerNotify to be much closer to BenchmarkUpdate than it is.

Other Benchmarks I Tried

A benchmark on a table with a trigger that doesn't do anything was 4% slower than BenchmarkUpdate. A trigger that inserts to a different table was 12% slower than BenchmarkUpdate. Nothing like the 5x slowdown a call to pg_notify does. So I ruled out the simple presence of the trigger to be the cause.

I noticed that updating the same row every time without any triggers runs way closer to BenchmarkUpdateNotifyTrigger performance than BenchmarkUpdate. That made me think that maybe pg_notify prevents Postgres from parallelizing the workload in some way, but I'm only guessing here.

I also tried making the trigger a CONSTRAINT trigger and adding DEFERRABLE INITIALLY DEFERRED parameters to it, to try to separate the transaction in which the update happens, and the notification is sent. It slowed BenchmarkUpdateNotifyTrigger further by ~22% instead of speeding it up.

How I Ran the Benchmarks

Put this code in a file main_test.go in an empty folder:

package main

import (
    "context"
    "sync"
    "testing"

    "github.com/jackc/pgx/v4"
)

const (
    concurrency = 10
    tableSize   = 100
    connString  = "postgres://test:test@localhost:5432/test"
)

func sqlInit(ctx context.Context, sql string) error {
    conn, err := pgx.Connect(ctx, connString)
    if err != nil {
        return err
    }
    defer conn.Close(ctx)

    _, err = conn.Exec(ctx, sql)
    return err
}

type sqlQuery struct {
    sql  string
    args []interface{}
}

type pool struct {
    sqlQueryCh chan sqlQuery
    wg         sync.WaitGroup

    err         error
    errMu       sync.Mutex
    nonNilErrCh chan struct{}
}

func newPool(ctx context.Context, cc int) (*pool, error) {
    ret := &pool{
        sqlQueryCh:  make(chan sqlQuery, cc),
        nonNilErrCh: make(chan struct{}),
    }

    for i := 0; i < cc; i++ {
        conn, err := pgx.Connect(ctx, connString)
        if err != nil {
            return nil, err
        }

        ret.wg.Add(1)
        go func(conn *pgx.Conn) {
            defer ret.wg.Done()
            defer conn.Close(ctx)

            for q := range ret.sqlQueryCh {
                if _, err := conn.Exec(ctx, q.sql, q.args...); err != nil {
                    ret.setErr(err)
                }
            }
        }(conn)
    }

    return ret, nil
}

func (p *pool) setErr(err error) {
    p.errMu.Lock()
    defer p.errMu.Unlock()
    if p.err == nil && err != nil {
        p.err = err
        close(p.nonNilErrCh)
    }
}

func (p *pool) Send(ctx context.Context, sql string, args ...interface{}) error {
    select {
    case p.sqlQueryCh <- sqlQuery{
        sql:  sql,
        args: args,
    }:
        return nil
    case <-ctx.Done():
        return ctx.Err()
    case <-p.nonNilErrCh:
        return p.err
    }
}

func (p *pool) Close() error {
    close(p.sqlQueryCh)
    p.wg.Wait()
    return p.err
}

func BenchmarkNotify(b *testing.B) {
    ctx := context.Background()

    if err := sqlInit(ctx, `
DROP TABLE IF EXISTS table1;
CREATE TABLE table1 (
    id SERIAL PRIMARY KEY,
    int INT
);`); err != nil {
        panic(err)
    }

    pool, err := newPool(ctx, concurrency)
    if err != nil {
        panic(err)
    }

    for i := 0; i < b.N; i++ {
        if err := pool.Send(ctx, "SELECT pg_notify('foo', 'bar');"); err != nil {
            panic(err)
        }
    }

    if err := pool.Close(); err != nil {
        panic(err)
    }
}

func BenchmarkUpdate(b *testing.B) {
    ctx := context.Background()

    if err := sqlInit(ctx, `
DROP TABLE IF EXISTS table1;
CREATE TABLE table1 (
    id SERIAL PRIMARY KEY,
    int INT
);
`); err != nil {
        panic(err)
    }

    pool, err := newPool(ctx, concurrency)
    if err != nil {
        panic(err)
    }

    for i := 0; i < b.N; i++ {
        if err := pool.Send(ctx, `
INSERT INTO table1 (id, int) VALUES($1, $2)
ON CONFLICT (id) DO UPDATE
SET int = $3`, i%tableSize, 2, 3); err != nil {
            panic(err)
        }
    }

    if err := pool.Close(); err != nil {
        panic(err)
    }
}

func BenchmarkUpdateTriggerNotify(b *testing.B) {
    ctx := context.Background()

    if err := sqlInit(ctx, `
DROP TABLE IF EXISTS table1;
CREATE TABLE table1 (
    id SERIAL PRIMARY KEY,
    int INT
);

CREATE OR REPLACE FUNCTION table1_fn() RETURNS TRIGGER AS $$
BEGIN
    PERFORM pg_notify('foo', 'bar');
    RETURN NEW;
END;
$$ LANGUAGE plpgsql;

DROP TRIGGER IF EXISTS table1_trg ON table1;
CREATE TRIGGER table1_trg AFTER UPDATE ON table1 FOR EACH ROW EXECUTE PROCEDURE table1_fn();
`); err != nil {
        panic(err)
    }

    pool, err := newPool(ctx, concurrency)
    if err != nil {
        panic(err)
    }

    for i := 0; i < b.N; i++ {
        if err := pool.Send(ctx, `
INSERT INTO table1 (id, int) VALUES($1, $2)
ON CONFLICT (id) DO UPDATE
SET int = $3`, i%tableSize, 2, 3); err != nil {
            panic(err)
        }
    }

    if err := pool.Close(); err != nil {
        panic(err)
    }
}

Run go mod init test.com/m

Run postgres docker image.

docker run -it -p 5433:5432 -e POSTGRES_USER=test -e POSTGRES_PASSWORD=test -e POSTGRES_DB=test postgres

Run the benchmarks in another window.

go test . --bench=.
goos: linux
goarch: amd64
pkg: example.com/m
cpu: Intel(R) Core(TM) i7-8550U CPU @ 1.80GHz
BenchmarkUpdate-8                       5182        309329 ns/op
BenchmarkUpdateTriggerNotify-8           895       1554109 ns/op
BenchmarkNotify-8                      27082         43838 ns/op
PASS
ok      example.com/m   7.840s
0 Answers
Related