> For the complete documentation index, see [llms.txt](https://docs.selfuel.digital/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://docs.selfuel.digital/data-integration-with-nexus/nexus-elements/connectors/sink/postgresql.md).

# PostgreSql

> JDBC PostgreSql Sink Connector

### Description[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#description) <a href="#description" id="description"></a>

Write data through jdbc. Support Batch mode and Streaming mode, support concurrent writing, support exactly-once semantics (using XA transaction guarantee).

### Key Features[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#key-features) <a href="#key-features" id="key-features"></a>

* [x] &#x20;exactly-once
* [x] &#x20;cdc

> Use `Xa transactions` to ensure `exactly-once`. So only support `exactly-once` for the database which is support `Xa transactions`. You can set `is_exactly_once=true` to enable it.

### Supported DataSource Info[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#supported-datasource-info) <a href="#supported-datasource-info" id="supported-datasource-info"></a>

<table><thead><tr><th>Datasource</th><th>Supported Versions</th><th>Driver</th><th>Url</th><th data-hidden>Maven</th></tr></thead><tbody><tr><td>PostgreSQL</td><td>Different dependency version has different driver class.</td><td>org.postgresql.Driver</td><td>jdbc:postgresql://localhost:5432/test</td><td><a href="https://mvnrepository.com/artifact/org.postgresql/postgresql">Download</a></td></tr><tr><td>PostgreSQL</td><td>If you want to manipulate the GEOMETRY type in PostgreSQL.</td><td>org.postgresql.Driver</td><td>jdbc:postgresql://localhost:5432/test</td><td><a href="https://mvnrepository.com/artifact/net.postgis/postgis-jdbc">Download</a></td></tr></tbody></table>

### Data Type Mapping[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#data-type-mapping) <a href="#data-type-mapping" id="data-type-mapping"></a>

| PostgreSQL Data Type                                                                            | Nexus Data Type                                                                                                                                |
| ----------------------------------------------------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------- |
| <p>BOOL<br></p>                                                                                 | BOOLEAN                                                                                                                                        |
| <p>\_BOOL<br></p>                                                                               | ARRAY\<BOOLEAN>                                                                                                                                |
| <p>BYTEA<br></p>                                                                                | BYTES                                                                                                                                          |
| <p>\_BYTEA<br></p>                                                                              | ARRAY\<TINYINT>                                                                                                                                |
| <p>INT2<br>SMALLSERIAL<br>INT4<br>SERIAL<br></p>                                                | INT                                                                                                                                            |
| <p>\_INT2<br>\_INT4<br></p>                                                                     | ARRAY\<INT>                                                                                                                                    |
| <p>INT8<br>BIGSERIAL<br></p>                                                                    | BIGINT                                                                                                                                         |
| <p>\_INT8<br></p>                                                                               | ARRAY\<BIGINT>                                                                                                                                 |
| <p>FLOAT4<br></p>                                                                               | FLOAT                                                                                                                                          |
| <p>\_FLOAT4<br></p>                                                                             | ARRAY\<FLOAT>                                                                                                                                  |
| <p>FLOAT8<br></p>                                                                               | DOUBLE                                                                                                                                         |
| <p>\_FLOAT8<br></p>                                                                             | ARRAY\<DOUBLE>                                                                                                                                 |
| NUMERIC(Get the designated column's specified column size>0)                                    | DECIMAL(Get the designated column's specified column size,Gets the number of digits in the specified column to the right of the decimal point) |
| NUMERIC(Get the designated column's specified column size<0)                                    | DECIMAL(38, 18)                                                                                                                                |
| <p>BPCHAR<br>CHARACTER<br>VARCHAR<br>TEXT<br>GEOMETRY<br>GEOGRAPHY<br>JSON<br>JSONB<br>UUID</p> | STRING                                                                                                                                         |
| <p>\_BPCHAR<br>\_CHARACTER<br>\_VARCHAR<br>\_TEXT</p>                                           | ARRAY\<STRING>                                                                                                                                 |
| <p>TIMESTAMP<br></p>                                                                            | TIMESTAMP                                                                                                                                      |
| <p>TIME<br></p>                                                                                 | TIME                                                                                                                                           |
| <p>DATE<br></p>                                                                                 | DATE                                                                                                                                           |
| OTHER DATA TYPES                                                                                | NOT SUPPORTED YET                                                                                                                              |

### Options[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#options) <a href="#options" id="options"></a>

| Name                                            | Type    | Required | Default                          | Description                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                        |
| ----------------------------------------------- | ------- | -------- | -------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| url                                             | String  | Yes      | -                                | <p>The URL of the JDBC connection. Refer to a case: jdbc:postgresql://localhost:5432/test<br>if you would use json or jsonb type insert please add jdbc url stringtype=unspecified option</p>                                                                                                                                                                                                                                                                                                                                                                                                                                                      |
| driver                                          | String  | Yes      | -                                | <p>The jdbc class name used to connect to the remote data source,<br>if you use PostgreSQL the value is <code>org.postgresql.Driver</code>.</p>                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                    |
| user                                            | String  | No       | -                                | Connection instance user name                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                      |
| password                                        | String  | No       | -                                | Connection instance password                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                       |
| query                                           | String  | No       | -                                | Use this sql write upstream input datas to database. e.g `INSERT ...`,`query` have the higher priority                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                             |
| database                                        | String  | No       | -                                | <p>Use this <code>database</code> and <code>table-name</code> auto-generate sql and receive upstream input datas write to database.<br>This option is mutually exclusive with <code>query</code> and has a higher priority.</p>                                                                                                                                                                                                                                                                                                                                                                                                                    |
| table                                           | String  | No       | -                                | <p>Use database and this table-name auto-generate sql and receive upstream input datas write to database.<br>This option is mutually exclusive with <code>query</code> and has a higher priority.The table parameter can fill in the name of an unwilling table, which will eventually be used as the table name of the creation table, and supports variables (<code>${table\_name}</code>, <code>${schema\_name}</code>). Replacement rules: <code>${schema\_name}</code> will replace the SCHEMA name passed to the target side, and <code>${table\_name}</code> will replace the name of the table passed to the table at the target side.</p> |
| primary\_keys                                   | Array   | No       | -                                | This option is used to support operations such as `insert`, `delete`, and `update` when automatically generate sql.                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                |
| support\_upsert\_by\_query\_primary\_key\_exist | Boolean | No       | false                            | Choose to use INSERT sql, UPDATE sql to process update events(INSERT, UPDATE\_AFTER) based on query primary key exists. This configuration is only used when database unsupport upsert syntax. **Note**: that this method has low performance                                                                                                                                                                                                                                                                                                                                                                                                      |
| connection\_check\_timeout\_sec                 | Int     | No       | 30                               | The time in seconds to wait for the database operation used to validate the connection to complete.                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                |
| max\_retries                                    | Int     | No       | 0                                | The number of retries to submit failed (executeBatch)                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                              |
| batch\_size                                     | Int     | No       | 1000                             | <p>For batch writing, when the number of buffered records reaches the number of <code>batch\_size</code> or the time reaches <code>checkpoint.interval</code><br>, the data will be flushed into the database</p>                                                                                                                                                                                                                                                                                                                                                                                                                                  |
| is\_exactly\_once                               | Boolean | No       | false                            | <p>Whether to enable exactly-once semantics, which will use Xa transactions. If on, you need to<br>set <code>xa\_data\_source\_class\_name</code>.</p>                                                                                                                                                                                                                                                                                                                                                                                                                                                                                             |
| generate\_sink\_sql                             | Boolean | No       | false                            | Generate sql statements based on the database table you want to write to.                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                          |
| xa\_data\_source\_class\_name                   | String  | No       | -                                | <p>The xa data source class name of the database Driver, for example, PostgreSQL is <code>org.postgresql.xa.PGXADataSource</code>, and<br>please refer to appendix for other data sources</p>                                                                                                                                                                                                                                                                                                                                                                                                                                                      |
| max\_commit\_attempts                           | Int     | No       | 3                                | The number of retries for transaction commit failures                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                              |
| transaction\_timeout\_sec                       | Int     | No       | -1                               | <p>The timeout after the transaction is opened, the default is -1 (never timeout). Note that setting the timeout may affect<br>exactly-once semantics</p>                                                                                                                                                                                                                                                                                                                                                                                                                                                                                          |
| auto\_commit                                    | Boolean | No       | true                             | Automatic transaction commit is enabled by default                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                 |
| field\_ide                                      | String  | No       | -                                | Identify whether the field needs to be converted when synchronizing from the source to the sink. `ORIGINAL` indicates no conversion is needed;`UPPERCASE` indicates conversion to uppercase;`LOWERCASE` indicates conversion to lowercase.                                                                                                                                                                                                                                                                                                                                                                                                         |
| properties                                      | Map     | No       | -                                | <p>Additional connection configuration parameters,when properties and URL have the same parameters, the priority is determined by the<br>specific implementation of the driver. For example, in MySQL, properties take precedence over the URL.</p>                                                                                                                                                                                                                                                                                                                                                                                                |
| common-options                                  |         | no       | -                                | Sink plugin common parameters, please refer to [Sink Common Options](/data-integration-with-nexus/nexus-elements/connectors/sink/sink-common-options.md) for details                                                                                                                                                                                                                                                                                                                                                                                                                                                                               |
| schema\_save\_mode                              | Enum    | no       | CREATE\_SCHEMA\_WHEN\_NOT\_EXIST | Before the synchronous task is turned on, different treatment schemes are selected for the existing surface structure of the target side.                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                          |
| data\_save\_mode                                | Enum    | no       | APPEND\_DATA                     | Before the synchronous task is turned on, different processing schemes are selected for data existing data on the target side.                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                     |
| custom\_sql                                     | String  | no       | -                                | When data\_save\_mode selects CUSTOM\_PROCESSING, you should fill in the CUSTOM\_SQL parameter. This parameter usually fills in a SQL that can be executed. SQL will be executed before synchronization tasks.                                                                                                                                                                                                                                                                                                                                                                                                                                     |
| enable\_upsert                                  | Boolean | No       | true                             | Enable upsert by primary\_keys exist, If the task has no key duplicate data, setting this parameter to `false` can speed up data import                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                            |

#### table \[string][​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#table-string) <a href="#table-string" id="table-string"></a>

Use `database` and this `table-name` auto-generate sql and receive upstream input datas write to database.

This option is mutually exclusive with `query` and has a higher priority.

The table parameter can fill in the name of an unwilling table, which will eventually be used as the table name of the creation table, and supports variables (`${table_name}`, `${schema_name}`). Replacement rules: `${schema_name}` will replace the SCHEMA name passed to the target side, and `${table_name}` will replace the name of the table passed to the table at the target side.

for example:

1. ${schema\_name}.${table\_name} \_test
2. dbo.tt\_${table\_name} \_sink
3. public.sink\_table

#### schema\_save\_mode\[Enum][​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#schema_save_modeenum) <a href="#schema_save_modeenum" id="schema_save_modeenum"></a>

Before the synchronous task is turned on, different treatment schemes are selected for the existing surface structure of the target side.\
Option introduction：\
`RECREATE_SCHEMA` ：Will create when the table does not exist, delete and rebuild when the table is saved\
`CREATE_SCHEMA_WHEN_NOT_EXIST` ：Will Created when the table does not exist, skipped when the table is saved\
`ERROR_WHEN_SCHEMA_NOT_EXIST` ：Error will be reported when the table does not exist

#### data\_save\_mode\[Enum][​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#data_save_modeenum) <a href="#data_save_modeenum" id="data_save_modeenum"></a>

Before the synchronous task is turned on, different processing schemes are selected for data existing data on the target side.\
Option introduction：\
`DROP_DATA`： Preserve database structure and delete data\
`APPEND_DATA`：Preserve database structure, preserve data\
`CUSTOM_PROCESSING`：User defined processing\
`ERROR_WHEN_DATA_EXISTS`：When there is data, an error is reported

#### custom\_sql\[String][​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#custom_sqlstring) <a href="#custom_sqlstring" id="custom_sqlstring"></a>

When data\_save\_mode selects CUSTOM\_PROCESSING, you should fill in the CUSTOM\_SQL parameter. This parameter usually fills in a SQL that can be executed. SQL will be executed before synchronization tasks.

#### Tips[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#tips) <a href="#tips" id="tips"></a>

> If partition\_column is not set, it will run in single concurrency, and if partition\_column is set, it will be executed in parallel according to the concurrency of tasks.

### Task Example[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#task-example) <a href="#task-example" id="task-example"></a>

#### Simple:[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#simple) <a href="#simple" id="simple"></a>

> This example defines a Nexus synchronization task that automatically generates data through FakeSource and sends it to JDBC Sink. FakeSource generates a total of 16 rows of data (row\.num=16), with each row having two fields, name (string type) and age (int type). The final target table is test\_table will also be 16 rows of data in the table. Before run this job, you need create database test and table test\_table in your PostgreSQL.&#x20;

```
# Defining the runtime environment
env {
  parallelism = 1
  job.mode = "BATCH"
}

source {
  FakeSource {
    parallelism = 1
    result_table_name = "fake"
    row.num = 16
    schema = {
      fields {
        name = "string"
        age = "int"
      }
    }
  }
  # If you would like to get more information about how to configure Nexus and see full list of source plugins,
  # please go to source page
}

transform {
  # If you would like to get more information about how to configure Nexus and see full list of transform plugins,
    # please go to transform page
}

sink {
    jdbc {
       # if you would use json or jsonb type insert please add jdbc url stringtype=unspecified option
        url = "jdbc:postgresql://localhost:5432/test"
        driver = "org.postgresql.Driver"
        user = root
        password = 123456
        query = "insert into test_table(name,age) values(?,?)"
     }
  # If you would like to get more information about how to configure Nexus and see full list of sink plugins,
  # please go to sink page
}
```

#### Generate Sink SQL[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#generate-sink-sql) <a href="#generate-sink-sql" id="generate-sink-sql"></a>

> This example not need to write complex sql statements, you can configure the database name table name to automatically generate add statements for you

```
sink {
    Jdbc {
        # if you would use json or jsonb type insert please add jdbc url stringtype=unspecified option
        url = "jdbc:postgresql://localhost:5432/test"
        driver = org.postgresql.Driver
        user = root
        password = 123456
        
        generate_sink_sql = true
        database = test
        table = "public.test_table"
    }
}
```

#### Exactly-once :[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#exactly-once-) <a href="#exactly-once" id="exactly-once"></a>

> For accurate write scene we guarantee accurate once

```
sink {
    jdbc {
       # if you would use json or jsonb type insert please add jdbc url stringtype=unspecified option
        url = "jdbc:postgresql://localhost:5432/test"
        driver = "org.postgresql.Driver"
    
        max_retries = 0
        user = root
        password = 123456
        query = "insert into test_table(name,age) values(?,?)"
    
        is_exactly_once = "true"
    
        xa_data_source_class_name = "org.postgresql.xa.PGXADataSource"
    }
}
```

#### CDC(Change Data Capture) Event[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#cdcchange-data-capture-event) <a href="#cdcchange-data-capture-event" id="cdcchange-data-capture-event"></a>

> CDC change data is also supported by us In this case, you need config database, table and primary\_keys.

```
sink {
    jdbc {
        # if you would use json or jsonb type insert please add jdbc url stringtype=unspecified option
        url = "jdbc:postgresql://localhost:5432/test"
        driver = "org.postgresql.Driver"
        user = root
        password = 123456
        
        generate_sink_sql = true
        # You need to configure both database and table
        database = test
        table = sink_table
        primary_keys = ["id","name"]
        field_ide = UPPERCASE
    }
}
```

#### Save mode function[​](https://seatunnel.apache.org/docs/2.3.7/connector-v2/sink/PostgreSql#save-mode-function) <a href="#save-mode-function" id="save-mode-function"></a>

```
sink {
    Jdbc {
        # if you would use json or jsonb type insert please add jdbc url stringtype=unspecified option
        url = "jdbc:postgresql://localhost:5432/test"
        driver = org.postgresql.Driver
        user = root
        password = 123456
        
        generate_sink_sql = true
        database = test
        table = "public.test_table"
        schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
        data_save_mode="APPEND_DATA"
    }
}
```
