Skip to main content

YugabyteDB continuous real-time replication using CDC

Yugabyte Reader uses the wal2json plugin to read change data. 1.x releases of wal2jon can not read transactions larger than 1 GB. If you have transactions that large, we recommend using a 2.x release of wal2json, which does not have that limitation. When you use a 2.x release, we recommend setting the Yugabyte Config property to automatically use format 2 for transactions larger than 1GB and format 1 for other transactions (see Yugabyte Reader properties).

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

If not using an automatic pipeline wizard, before creating the continuous real-time replication application, see Switching from initial load to continuous replication of YugabyteDB sources.

Configuring YugabyteDB to use Yugabyte Reader

Striim reads change data from YugabyteDB.

Yugabyte Reader requires logical replication. For general information about logical replication, see YugabyteDB / Explore / Change data capture / PostgreSQL protocol / Key concepts / PostgreSQL logical replication concepts. To capture CDC changes, enter the following in a YugabyteDB client to create a replication slot:

SELECT pg_create_logical_replication_slot('SLOT_NAME', 'wal2json');

Enter the following command to create a role (user) for use by Striim, replacing ****** with a strong password:

CREATE ROLE striim WITH LOGIN PASSWORD '******' REPLICATION;

Yugabyte setup for schema evolution

Using Schema evolution with Yugabyte Reader requires a tracking table in the source database. To create this table, run the following commands in a YugabyteDB client:

CREATE SCHEMA IF NOT EXISTS striim;
CREATE TABLE IF NOT EXISTS striim.ddlcapturetable
  (
    event           TEXT,
    tag             TEXT,
    classid         OID,
    objid           OID,
    objsubid        INT,
    object_type     TEXT,
    schema_name     TEXT,
    object_identity TEXT,
    is_extension    BOOL,
    query           TEXT,
    username        TEXT DEFAULT CURRENT_USER,
    db_name TEXT DEFAULT Current_database(),
    client_addr     INET DEFAULT Inet_client_addr(),
    creation_time   TIMESTAMP DEFAULT now(),
    id SERIAL PRIMARY KEY
  ); 
GRANT USAGE ON SCHEMA striim TO PUBLIC;
GRANT SELECT, INSERT ON TABLE striim.ddlcapturetable TO PUBLIC;

Yugabyte programmer's reference

Yugabyte Reader properties

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

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

property

type

default value

notes

Bidirectional Marker Table

String

Not supported in this release.

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.

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.

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

Replication Slot Name

String

striim_slot

The name of the replication slot created as described in Configuring YugabyteDB to use Yugabyte Reader. If you have multiple instances of Yugabyte Reader, 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

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.

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 Yugabyte Reader, each should read a separate set of tables.

Username

String

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

Yugabyte Config

String

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

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

Target data type support & mapping for YugabyteDB sources

The table below details how Striim maps the data types of a YugabyteDB 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.

YugabyteDB 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.