Skip to main content
Version: 3.0.0

Databend

Databend sink connector

Support Those Engines​

Spark
Flink
SeaTunnel Zeta

Key Features​

Description​

A sink connector for writing data to Databend. Supports both batch and streaming processing modes. The Databend sink internally implements bulk data import through stage attachment.

Dependencies​

  1. You need to download the Databend JDBC driver jar package and add it to the directory ${SEATUNNEL_HOME}/plugins/.

For SeaTunnel Zeta​

  1. You need to download the Databend JDBC driver jar package and add it to the directory ${SEATUNNEL_HOME}/lib/.

Supported DataSource Info​

In order to use the Databend connector, the following dependencies are required. They can be downloaded via install-plugin.sh or from the Maven central repository.

DatasourceSupported VersionsDependency
Databend1.2.x and aboveDownload

Sink Options​

NameTypeRequiredDefault ValueDescription
urlStringYes-Databend JDBC connection URL. It must start with jdbc:databend://
usernameStringYes-Databend database username
passwordStringYes-Databend database password
databaseStringNo-Databend database name, defaults to the database name specified in the connection URL
tableStringNo-Databend table name
batch_sizeIntegerNo1000Number of records for batch writing
auto_commitBooleanNotrueWhether to auto-commit transactions
max_retriesIntegerNo3Maximum retry attempts on write failure
schema_save_modeEnumNoCREATE_SCHEMA_WHEN_NOT_EXISTSchema save mode
data_save_modeEnumNoAPPEND_DATAData save mode
custom_sqlStringNo-Custom write SQL, typically used for complex write scenarios
execute_timeout_secIntegerNo300SQL execution timeout (seconds)
jdbc_configMapNo-Additional JDBC connection configuration, such as connection timeout parameters
conflict_keyStringNo-Conflict key for CDC mode, used to determine the primary key for conflict resolution
enable_deleteBooleanNofalseWhether to allow delete operations in CDC mode
common-optionsNo-Sink plugin common parameters, please refer to Sink Common Options for details.

schema_save_mode [Enum]​

Before starting the synchronization task, choose different processing schemes for existing table structures. Option descriptions:
RECREATE_SCHEMA: Create when table doesn't exist, drop and recreate when table exists.
CREATE_SCHEMA_WHEN_NOT_EXIST: Create when table doesn't exist, skip when table exists.
ERROR_WHEN_SCHEMA_NOT_EXIST: Report error when table doesn't exist.
IGNORE: Ignore table processing.

data_save_mode [Enum]​

Before starting the synchronization task, choose different processing schemes for existing data on the target side. Option descriptions:
DROP_DATA: Retain database structure and delete data.
APPEND_DATA: Retain database structure and data.
CUSTOM_PROCESSING: User-defined processing.
ERROR_WHEN_DATA_EXISTS: Report error when data exists.

Data Type Mapping​

SeaTunnel Data TypeDatabend Data Type
BOOLEANBOOLEAN
TINYINTTINYINT
SMALLINTSMALLINT
INTINT
BIGINTBIGINT
FLOATFLOAT
DOUBLEDOUBLE
DECIMALDECIMAL
STRINGSTRING
BYTESVARBINARY
DATEDATE
TIMETIME
TIMESTAMPTIMESTAMP

Task Examples​

Simple Example​

env {
execution.parallelism = 1
job.mode = "BATCH"
}

source {
FakeSource {
row.num = 10
schema = {
fields {
name = string
age = int
score = double
}
}
}
}

sink {
Databend {
url = "jdbc:databend://localhost:8000"
username = "root"
password = ""
database = "default"
table = "target_table"
batch_size = 1000
}
}

Writing with Custom SQL​

sink {
Databend {
url = "jdbc:databend://localhost:8000"
username = "root"
password = ""
database = "default"
table = "target_table"
custom_sql = "INSERT INTO default.target_table(name, age, score) VALUES(?, ?, ?)"
}
}

Using Schema Save Mode​

sink {
Databend {
url = "jdbc:databend://localhost:8000"
username = "root"
password = ""
database = "default"
table = "target_table"
schema_save_mode = "RECREATE_SCHEMA"
data_save_mode = "APPEND_DATA"
}
}

CDC mode​

Set conflict_key to the primary-key column used to merge update/delete events. Set enable_delete = true only when DELETE events should remove rows from Databend. If conflict_key is not configured, the sink writes normal insert-style batches.

The following end-to-end example feeds CDC row kinds (INSERT, UPDATE_BEFORE, UPDATE_AFTER, DELETE) into Databend. The sink merges updates and applies deletes against the rows identified by conflict_key.

env {
parallelism = 1
job.mode = "BATCH"
checkpoint.interval = 1000
}

source {
FakeSource {
row.num = 10
schema = {
fields {
id = "int"
name = "string"
position = "string"
age = "int"
score = "double"
}
}
rows = [
{
kind = INSERT
fields = [1, "Alice", "Engineer", 30, 95.5]
},
{
kind = INSERT
fields = [2, "Bob", "Developer", 25, 85.0]
},
{
kind = UPDATE_BEFORE
fields = [2, "Bob", "Developer", 25, 85.0]
},
{
kind = UPDATE_AFTER
fields = [2, "Bob", "Senior Developer", 25, 87.0]
},
{
kind = DELETE
fields = [2, "Bob", "Senior Developer", 25, 87.0]
}
]
}
}

sink {
Databend {
url = "jdbc:databend://databend:8000/default?ssl=false"
username = "root"
password = ""
database = "default"
table = "sink_table"

# Enable CDC mode
batch_size = 1
conflict_key = "id"
enable_delete = true
}
}

Stream MySQL CDC To Databend In Streaming Mode​

The same CDC settings also work in streaming jobs. The following example pipes MySQL CDC events into Databend continuously. Keep batch_size small in streaming CDC jobs so that each checkpoint reflects the latest writes:

env {
parallelism = 1
job.mode = "STREAMING"
checkpoint.interval = 10000
}

source {
MySQL-CDC {
base-url = "jdbc:mysql://mysql:3306/test"
username = "root"
password = "mysqlpw"
table-names = ["test.orders"]
}
}

sink {
Databend {
url = "jdbc:databend://databend:8000/default?ssl=false"
username = "root"
password = ""
database = "default"
table = "orders"
batch_size = 500
conflict_key = "id"
enable_delete = true
}
}

Changelog​

Change Log
ChangeCommitVersion
[Fix][CDC][Zeta] Restore runtime schema from checkpoint after failover (#11503)https://github.com/apache/seatunnel/commit/ec1b1b8b53.0.0
[Fix][Connector-V2] Fix CDC comment schema-change event routing (#11837)https://github.com/apache/seatunnel/commit/ed8d941513.0.0
[Improve][Connector-V2] Validate Databend JDBC URL (#11831)https://github.com/apache/seatunnel/commit/4a65e98ad3.0.0
[Feature][Connector-V2][CDC] Support comment-related schema change events (#11025)https://github.com/apache/seatunnel/commit/ba55ef9653.0.0
[Feature][Connector-V2] Support databend source/sink connector (#9331)https://github.com/apache/seatunnel/commit/2f96f2e46c2.3.12