Skip to main content

PostgreSQL programmer's reference

PostgreSQL Reader properties

Before you can use this adapter, PostgreSQL must be configured as described in Configuring PostgreSQL to use PostgreSQL Reader.PostgreSQL setup

Striim provides wizards for creating applications that read from PostgreSQL and write to various targets. See Creating an application using a wizard for details.

property

type

default value

notes

Bidirectional Marker Table

String

CDDL Action

enum

Process

Appears in Flow Designer only when CDDL Capture is True. See Handling schema evolution.

CDDL Capture

Boolean

False

See Handling schema evolution.

Do not use Find and Replace DDL unless instructed to by Striim support.

CDDL Tracking Table

String

Appears in Flow Designer only when CDDL Capture is True. See PostgreSQL setup for schema evolution.

Connection Profile Name

String

Appears in Flow Designer only when Use Connection Profile is True. See Connection profiles.

Connection Retry Policy

String

retryInterval=30, maxRetries=3

With the default setting, if a connection attempt is unsuccessful, the adapter will try again in 30 seconds (retryInterval. If the second attempt is unsuccessful, in 30 seconds it will try a third time (maxRetries). If that is unsuccessful, the adapter will fail and log an exception. Negative values are not supported.

Connection URL

String

jdbc:postgresql:// followed by the primary server's IP address or network name, a colon, the port number, and a slash followed by the database name. If the database name is omitted, the Username value is used as the database name.

PostgreSQL Reader cannot read from a replica (standby) server since the replication slot is in the primary server.

Excluded Tables

String

Data for any tables specified here will not be returned. For example, if Tables uses a wildcard, data from any tables specified here will be omitted. Multiple table names (separated by semicolons) and wildcards may be used exactly as for Tables.

Filter Transaction Boundaries

Boolean

True

With the default value of True, begin and commit transactions are filtered out. Set to False to include begin and commit transactions. This must be set to False to enable Preserve Source Transaction Boundary in a downstream writer.

Password

encrypted password

the password specified for the username (see Encrypted passwords)

Postgres Config

String

{"ReplicationPluginConfig": {"Name": "WAL2JSON", "Format": "1"}}

Change 1 to 2 to use wal2json format 2 (see the wal2json readme for more information).

If you are running an older version of Amazon RDS for PostgreSQL that supports only version 1, you may contact AWS technical support to have the wal2json plugin updated.

Replication Slot Name

String

striim_slot

The name of the replication slot created as described in Configuring PostgreSQL to use PostgreSQL Reader. If you have multiple instances of PostgreSQLReader, each must have its own slot.

Start LSN

String

By default, only new transactions are read. Optionally, specify a log sequence number to start reading from that point.

If you are using schema evolution (see Handling schema evolution, set a Start LSN only if you are sure that there have been no DDL changes after that point.Handling schema evolution

See Switching from initial load to continuous replication of PostgreSQL sources.

Tables

String

The table(s) for which to return change data. Tables must have primary keys or REPLICA IDENTITY set to FULL (required for logical replication).

Names are case-sensitive. Specify source table names as <schema>.<table>) (The database is specified in the connection URL.)

Do not modify this property when CDDL Capture is True or recovery is enabled for the application.

You may specify multiple tables as a list separated by semicolons or using the following wildcards in the schema and/or table names only (not in the database name):

  • %: any series of characters

  • _: any single character

For example, %.% would include all tables in all schemas in the database specified in the connection URL.

The % wildcard is allowed only at the end of the string. For example, mydb.prefix% is valid, but mydb.%suffix is not. If CDDL Capture is enabled, partial wildcards such as mydb.prefix% will not pick up new tables added after the application is started. To work around this limitation, use a full wildcard, such as mydb.%.

All tables specified must have primary keys. Tables without primary keys are not included in output.

Known issue DEV-27752: this adapter can not read partitioned tables. See wal2json issue #259.

If any specified tables are missing Striim will issue a warning. If none of the specified tables exists, start will fail with a "found no tables" error.

If you have multiple instances of PostgreSQLReader, each should read a separate set of tables.

Use Connection Profile

Boolean

False

Set to True to use a connection profile instead of specifying the connection properties in the adapter properties. See Connection profiles.

Username

String

the login name for the user created as described in Configuring PostgreSQL to use PostgreSQL Reader

PostgreSQL Reader WAEvent fields

The output data type for PostgreSQLReader is WAEvent. The elements are:

metadata: a map including:

  • LSN: log sequence number of the transaction's commit

  • NEXT_LSN: next log sequence number (used for reconnecting to the replication slot after a non-fatal network interruption)

  • OperationName: INSERT, UPDATE, or DELETE

    When schema evolution is enabled, OperationName for DDL events will be Alter, AlterColumns, Create, or Drop. This metadata is reserved for internal use by Striim and subject to change, so should not be used in CQs, open processors, or custom Java functions.

  • PK_UPDATE: included only when an UPDATE changes the primary key

  • Sequence: incremented for each operation within a transaction

  • TableName: the name of the table including its schema

  • Timestamp: timestamp from the replication subscription

  • TxnID: transaction identifier

To retrieve the values for these fields, use the META() function. See Parsing the fields of WAEvent for CDC readers.

data: an array of fields, numbered from 0, containing:

  • for an INSERT operation, the values that were inserted

  • for an UPDATE, the values after the operation was completed

  • for a DELETE, the value of the primary key and nulls for the other fields

To retrieve the values for these fields, use SELECT ... (DATA[]). See Parsing the fields of WAEvent for CDC readers.

before: for UPDATE operations, contains the primary key value from before the update. When an update changes the primary key value, you may retrieve the previous value using the BEFORE() function.

dataPresenceBitMap, beforePresenceBitMap, and typeUUID are reserved and should be ignored.

PostgreSQL Reader simple application

The following application will write change data for all tables in all schemas in database mydb to SysOut. Replace striim and ****** with the user name and password for the PostgreSQL account you created for use by PostgreSQLReader (see Configuring PostgreSQL to use PostgreSQL Reader) and mydb and %.% with the names of the database and tables to be read. If the replication slot name is not striim_slot, specify it using the ReplicationSlotName property.

CREATE APPLICATION PostgreSQLTest;

CREATE SOURCE PostgreSQLCDCIn USING PostgreSQLReader (
  Username:'striim',
  Password:'******',
  ConnectionURL:'jdbc:postgresql://192.0.2.10:5432/mydb',
  ReplicationSlotName: 'striim_slot',
  Tables:'%.%'
) 
OUTPUT TO PostgreSQLCDCStream;

CREATE TARGET PostgreSQLCDCOut
USING SysOut(name:PostgreSQLCDC)
INPUT FROM PostgreSQLCDCStream;

END APPLICATION PostgreSQLTest;

PostgreSQL Reader example output

PostgreSQLReader's output type is WAEvent. See WAEvent contents for change data and PostgreSQL Reader WAEvent fields for more information.

The following are examples of WAEvents emitted by PostgreSQLReader for various operation types. They all use the following table:

CREATE TABLE posauthorizations (
  business_name varchar(30),
  merchant_id character varying(35) PRIMARY KEY,
  primary_account bigint,
  pos bigint,
  code character varying(20),
  exp character(4),
  currency_code character(3),
  auth_amount numeric(10,3),
  terminal_id bigint,
  zip bigint,
  city character varying(20));
INSERT

If you performed the following INSERT on the table:

INSERT INTO posauthorizations VALUES(
  'COMPANY 1',
  'D6RJPwyuLXoLqQRQcOcouJ26KGxJSf6hgbu',
  6705362103919221351,
  0,
  '20130309113025',
  '0916',
  'USD',
  2.20,
  5150279519809946,
  41363,
  'Quicksand');

The WAEvent for that INSERT would be similar to:

data: ["COMPANY 1","D6RJPwyuLXoLqQRQcOcouJ26KGxJSf6hgbu",6705362103919221351,0,"20130309113025",
"0916","USD",2.200,5150279519809946,41363,"Quicksand"]
metadata: {"TableName":"public.posauthorizations","TxnID":556,"OperationName":"INSERT",
"LSN":"0/152CD58","NEXT_LSN":"0/152D1C8","Sequence":1,"Timestamp":"2019-01-11 16:29:54.628403-08"}
UPDATE

If you performed the following UPDATE on the table:

UPDATE posauthorizations SET BUSINESS_NAME = 'COMPANY 5A' where pos=0;

The WAEvent for that UPDATE would be similar to:

data: ["COMPANY 5A","D6RJPwyuLXoLqQRQcOcouJ26KGxJSf6hgbu",6705362103919221351,0,"20130309113025",
"0916","USD",2.200,5150279519809946,41363,"Quicksand"]
metadata: {"TableName":"public.posauthorizations","TxnID":557,"OperationName":"UPDATE",
"LSN":"0/152D2E0","NEXT_LSN":"0/152D6F8","Sequence":1,"Timestamp":"2019-01-11 16:31:54.271525-08"}
before: [null,"D6RJPwyuLXoLqQRQcOcouJ26KGxJSf6hgbu",null,null,null,null,null,null,null,null,null]

When an UPDATE changes the primary key, you may retrieve the old primary key value from the before array.

DELETE

If you performed the following DELETE on the table:

DELETE from posauthorizations where pos=0;

The WAEvent for that DELETE would be similar to:

data: [null,"D6RJPwyuLXoLqQRQcOcouJ26KGxJSf6hgbu",null,null,null,null,null,null,null,null,null]
metadata: {"TableName":"public.posauthorizations","TxnID":558,"OperationName":"DELETE",
"LSN":"0/152D730","NEXT_LSN":"0/152D7C8","Sequence":1,"Timestamp":"2019-01-11 16:33:09.065951-08"}

Only the primary key value is included.

PostgreSQL Reader data type support and correspondence

Data types created using CREATE TYPE are not supported.

PostgreSQL type

Striim type

bigint

long

bigserial

long

bit

string

bit varying

string

boolean

short

bytea

string

character

string

character varying

string

cidr

string

circle

unsupported

composite type

string

date

DateTime

daterange

string

double precision

double

inet

string

integer

integer

int2

short

int4

integer

int4range

string

int8

long

int8range

string

integer

integer

interval

string

json

string

jsonb

string

line

unsupported

lseg

unsupported

macaddr

string

macaddr8

string

money

string

name (system identifier)

string

numeric

string (Infinity, -Infinity, and NaN values will be converted to null)

numrange

string

path

unsupported

pg_lan

string

point

unsupported

polygon

unsupported

real

float

smallint

short

smallserial

short

serial

integer

text

string

time

string

time with time zone

string

timestamp

datetime

tsrange

string

timestamp with time zone

datetime

tstzrange

string

tsquery

unsupported

tsvector

unsupported

txid_snapshot

string

uuid

string

xml

string

Target data type support & mapping for PostgreSQL sources

The table below details how Striim maps the data types of a PostgreSQL source to ClickHouse data types when you create an application using a wizard with Auto Schema Creation, perform an initial load using Database Reader with Create Schema enabled, or run the schema conversion utility, or when Striim schema evolution creates or alters target tables.

See Data Types for a list of supported data type aliases (such as decimal and varchar).

For fixed-length data types, Striim interprets the length parameter as one character = one byte, which can result in errors if the data uses multi-byte characters. To avoid this issue, manually increase the size of the data type in the target, or change the target data type to blob or clob.

PostgreSQL data types created using CREATE TYPE are not supported.

PostgreSQL Data Type

ClickHouse

BIGSERIAL

Int64

BIT

FixedString(1)

BIT(p)

String, if (p) > 1*

BOOL

Bool

BOX

String

BPCHAR

String

BPCHAR(p)

String

BYTEA

String

CIDR

String

CIRCLE

String

DATE

Date32

DATERANGE

String

FLOAT4

Float32

FLOAT8

Float64

INET

String

INT2

Int16

INT4

Int32

INT4RANGE

String

INT8

Int64

INT8RANGE

String

INTERVAL

String

INTERVAL(p)

String

JSON

String

JSONB

String

LINE

String

LSEG

String

MACADDR

String

MONEY

String

NUMERIC

Decimal(76)

NUMERIC(p,0)

Decimal(p, s), if (p) <= 76, if (s) <= 76

NUMERIC(p,s)

Decimal(p, s), if (p) <= 76, if (s) <= 76

String, if (p,s) > 76, if (s) > 76*

NUMRANGE

String

PATH

String

POINT

String

POLYGON

String

SERIAL

Int32

SMALLSERIAL

Int16

TEXT

String

TIME

Time64(s)

TIME(p)

Time64(s), if (s) <= 9

TIMESTAMP

DateTime64(s)

TIMESTAMP(p)

DateTime64(s), if (s) <= 9

TIMESTAMPTZ

DateTime64(s)

TIMESTAMPTZ(p)

DateTime64(s), if (s) <= 9

TIMETZ

String

TIMETZ(p)

String

TSQUERY

String

TSRANGE

String

TSTZRANGE

String

TSVECTOR

String

TXID_SNAPSHOT

String

UUID

UUID

VARBIT

String

VARBIT(p)

String

VARCHAR

String

VARCHAR(p)

String

XML

String

*When using the schema conversion utility, these mappings appear in converted_tables_with_striim_intelligence.sql.