Skip to main content
Version: Next

Tablestore

Tablestore sink connector

Support Those Engines

Spark
Flink
SeaTunnel Zeta

Description

Write SeaTunnel rows to Alibaba Cloud Tablestore.

Key features

Data Type Mapping

SeaTunnel typeTablestore attribute column typeTablestore primary key type
INT, TINYINT, SMALLINT, BIGINTINTEGERINTEGER
FLOAT, DOUBLE, DECIMALDOUBLESTRING
STRING, DATE, TIME, TIMESTAMPSTRINGSTRING
BOOLEANBOOLEANSTRING
BYTESBINARYBINARY

Options

nametyperequireddefault valuedescription
end_pointstringyes-Tablestore endpoint.
instance_namestringyes-Tablestore instance name.
access_key_idstringyes-AccessKey ID used to access Tablestore.
access_key_secretstringyes-AccessKey secret used to access Tablestore.
tablestringyes-Target Tablestore table name.
primary_keysarrayyes-Primary key field names in the target table.
schemaconfigyes-Input schema. Primary key fields must also exist in schema.fields.
batch_sizeintno25Maximum number of rows written in one batch.
common-optionsconfigno-Sink common options.

Usage notes

  • primary_keys can contain one or more primary key fields. These fields are written as Tablestore primary key columns; all other schema fields are written as normal attribute columns.
  • The sink writes rows with Tablestore RowPutChange and RowExistenceExpectation.IGNORE. It does not delete rows when upstream sends DELETE row kinds.
  • batch_size controls when buffered rows are flushed. The writer also flushes remaining rows when the job closes.
  • Keep access_key_id and access_key_secret out of committed job files. Prefer runtime variable substitution or a secret manager supported by your deployment environment.

common options [config]

Sink plugin common parameters, please refer to Sink Common Options for details.

Example

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

source {
FakeSource {
row.num = 2
schema = {
fields {
order_id = string
user_id = string
amount = double
}
}
rows = [
{
fields = ["order-1", "user-1", 99.5]
},
{
fields = ["order-2", "user-2", 20.0]
}
]
}
}

sink {
Tablestore {
end_point = "https://<instance>.<region>.ots.aliyuncs.com"
instance_name = "<instance-name>"
access_key_id = "${ACCESS_KEY_ID}"
access_key_secret = "${ACCESS_KEY_SECRET}"
table = "orders"
primary_keys = ["order_id"]
batch_size = 25
schema = {
fields {
order_id = string
user_id = string
amount = double
}
}
}
}

Changelog

Change Log
ChangeCommitVersion
[Improve] table_store options (#9515)https://github.com/apache/seatunnel/commit/145b68793f2.3.12
[Feature][Checkpoint] Add check script for source/sink state class serialVersionUID missing (#9118)https://github.com/apache/seatunnel/commit/4f5adeb1c72.3.11
[Improve] restruct connector common options (#8634)https://github.com/apache/seatunnel/commit/f3499a6eeb2.3.10
[Feature][Restapi] Allow metrics information to be associated to logical plan nodes (#7786)https://github.com/apache/seatunnel/commit/6b7c53d03c2.3.9
[Feature][Connector-V2][Tablestore] Support Source connector for Tablestore #7448 (#7467)https://github.com/apache/seatunnel/commit/a7ca51b5852.3.8
[Improve][Common] Introduce new error define rule (#5793)https://github.com/apache/seatunnel/commit/9d1b2582b22.3.4
[Improve] Remove use SeaTunnelSink::getConsumedType method and mark it as deprecated (#5755)https://github.com/apache/seatunnel/commit/8de74081002.3.4
Support config column/primaryKey/constraintKey in schema (#5564)https://github.com/apache/seatunnel/commit/eac76b4e502.3.4
[Improve][Connector-V2] Remove scheduler in Tablestore sink (#5272)https://github.com/apache/seatunnel/commit/8d6b07e4662.3.3
Merge branch 'dev' into merge/cdchttps://github.com/apache/seatunnel/commit/4324ee19122.3.1
[Improve][Project] Code format with spotless plugin.https://github.com/apache/seatunnel/commit/423b5830382.3.1
[improve][api] Refactoring schema parse (#4157)https://github.com/apache/seatunnel/commit/b2f573a13e2.3.1
[Improve][build] Give the maven module a human readable name (#4114)https://github.com/apache/seatunnel/commit/d7cd6010512.3.1
[Improve][Project] Code format with spotless plugin. (#4101)https://github.com/apache/seatunnel/commit/a2ab1665612.3.1
[Hotfix][OptionRule] Fix option rule about all connectors (#3592)https://github.com/apache/seatunnel/commit/226dc6a1192.3.0
[Improve][Connector-V2][TableStore] Unified excetion for TableStore sink connector (#3527)https://github.com/apache/seatunnel/commit/7b264d70042.3.0
[Feature][connector-v2] add tablestore source and sink (#3309)https://github.com/apache/seatunnel/commit/ebebf0b6332.3.0