跳到主要内容
版本:3.0.0

Doris

Doris 源连接器

支持的引擎​

Spark
Flink
SeaTunnel Zeta

主要功能​

描述​

用于 Apache Doris 的源连接器。

依赖​

  1. 你需要下载 jdbc driver jar package 并添加到目录 ${SEATUNNEL_HOME}/plugins/.

对于 SeaTunnel Zeta​

  1. 你需要下载 jdbc driver jar package 并添加到目录 ${SEATUNNEL_HOME}/lib/.

支持的数据源信息​

数据源支持版本驱动UrlMaven
Doris仅支持Doris2.0及以上版本.---

数据类型映射​

Doris 数据类型SeaTunnel 数据类型
INTINT
TINYINTTINYINT
SMALLINTSMALLINT
BIGINTBIGINT
LARGEINTSTRING
BOOLEANBOOLEAN
DECIMALDECIMAL((Get the designated column's specified column size)+1,
(Gets the designated column's number of digits to right of the decimal point.)))
FLOATFLOAT
DOUBLEDOUBLE
CHAR
VARCHAR
STRING
TEXT
STRING
JSONSTRING
VARIANTSTRING
DATEDATE
DATETIME
DATETIME(p)
TIMESTAMP
ARRAYARRAY

源选项​

基础配置:

名称类型是否必须默认值描述
fenodesstringyes-FE 地址, 格式:"fe_host:fe_http_port"
usernamestringyes-用户名
passwordstringyes-密码
doris.request.retriesintno3请求Doris FE的重试次数
doris.request.read.timeout.msintno30000请求 Doris BE 的 socket 读取超时时间。
doris.request.connect.timeout.msintno30000请求 Doris FE 或 BE 的连接超时时间。
query-portintno9030Doris 查询端口。
doris.request.query.timeout.sintno3600Doris扫描数据的超时时间,单位秒
doris.request.tablet.sizeintnoInteger.MAX_VALUE每个 SeaTunnel split 包含的 Doris tablet 数量,最小值为 1。
doris.deserialize.arrow.asyncbooleannofalse是否异步反序列化 Arrow 数据。
doris.request.retriesdoris.deserialize.queue.sizeintno64异步反序列化 Arrow 数据时使用的队列大小。
table_listArrayno-要读取的 Doris 表清单。

doris.request.retriesdoris.deserialize.queue.size 是当前运行时实际使用的配置名。调整异步 Arrow 反序列化队列大小时,请按这个完整名称配置。

表清单配置:

名称类型是否必须默认值描述
databasestringyes-数据库
tablestringyes-表名
doris.read.fieldstringno-选择要读取的Doris表字段
doris.filter.querystringno-数据过滤. 格式:"字段 = 值", 例如:doris.filter.query = "F_ID > 2"
doris.request.tablet.sizeintnoInteger.MAX_VALUE当前表每个 SeaTunnel split 包含的 Doris tablet 数量,最小值为 1。
doris.batch.sizeintno1024每次能够从BE中读取到的最大行数
doris.exec.mem.limitlongno2147483648单个be扫描请求可以使用的最大内存。默认内存为2G(2147483648)

注意: 当此配置对应于单个表时,您可以将table_list中的配置项展平到外层。如果不配置 table_list,必须在 source 外层配置 database 和 table。

提示​

不建议随意修改高级参数

例子​

单表​

这是一个从doris读取数据后,输出到控制台的例子:

env {
parallelism = 2
job.mode = "BATCH"
}
source{
Doris {
fenodes = "doris_e2e:8030"
username = root
password = ""
database = "e2e_source"
table = "doris_e2e_table"
}
}

transform {
# If you would like to get more information about how to configure seatunnel and see full list of transform plugins,
# please go to https://seatunnel.apache.org/docs/transforms/sql
}

sink {
Console {}
}

使用doris.read.field参数来选择需要读取的Doris表字段:

env {
parallelism = 2
job.mode = "BATCH"
}
source{
Doris {
fenodes = "doris_e2e:8030"
username = root
password = ""
database = "e2e_source"
table = "doris_e2e_table"
doris.read.field = "F_ID,F_INT,F_BIGINT,F_TINYINT,F_SMALLINT"
}
}

transform {
# If you would like to get more information about how to configure seatunnel and see full list of transform plugins,
# please go to https://seatunnel.apache.org/docs/transforms/sql
}

sink {
Console {}
}

使用doris.filter.query来过滤数据,参数值将作为过滤条件直接传递到doris:

env {
parallelism = 2
job.mode = "BATCH"
}
source{
Doris {
fenodes = "doris_e2e:8030"
username = root
password = ""
database = "e2e_source"
table = "doris_e2e_table"
doris.filter.query = "F_ID > 2"
}
}

transform {
# If you would like to get more information about how to configure seatunnel and see full list of transform plugins,
# please go to https://seatunnel.apache.org/docs/transforms/sql
}

sink {
Console {}
}

多表​

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

source{
Doris {
fenodes = "xxxx:8030"
username = root
password = ""
table_list = [
{
database = "st_source_0"
table = "doris_table_0"
doris.read.field = "F_ID,F_INT,F_BIGINT,F_TINYINT"
doris.filter.query = "F_ID >= 50"
doris.request.tablet.size = 1
doris.exec.mem.limit = 2147483648
},
{
database = "st_source_1"
table = "doris_table_1"
}
]
}
}

transform {}

sink{
Doris {
fenodes = "xxxx:8030"
schema_save_mode = "RECREATE_SCHEMA"
username = root
password = ""
database = "st_sink"
table = "${table_name}"
sink.enable-2pc = "true"
sink.label-prefix = "test_json"
doris.config = {
format="json"
read_json_by_line="true"
}
}
}

变更日志​

Change Log
ChangeCommitVersion
[Fix][CDC][Zeta] Restore runtime schema from checkpoint after failover (#11503)https://github.com/apache/seatunnel/commit/ec1b1b8b53.0.0
[Improve][Connector-V2][Doris] Support partition cleanup for DROP_DATA (#11917)https://github.com/apache/seatunnel/commit/2714e6e253.0.0
[Improve][Connector-V2] Remove redundant imperative validation in Doris connector (#11858)https://github.com/apache/seatunnel/commit/d8186ba2d3.0.0
[Feature][shade]Refactor the seatunnel-shade module. (#9993)https://github.com/apache/seatunnel/commit/4ba2895953.0.0
[Feature][Connector-V2][CDC] Support comment-related schema change events (#11025)https://github.com/apache/seatunnel/commit/ba55ef9653.0.0
[Fix][Connector-V2] Cap decimal scale to what Doris 1.x accepts (#11690)https://github.com/apache/seatunnel/commit/1cbff54fe3.0.0
[Feature][Connector-V2] Support timer flush for Doris sink (#11506)https://github.com/apache/seatunnel/commit/3ba5470a33.0.0
[Fix][Connector-V2] Fix Doris stream load waitForContinue timeout when FE redirect is slow (#11330)https://github.com/apache/seatunnel/commit/e6c0973913.0.0
[Feature][Doris] Support VARIANT source and sink mapping (#10855)https://github.com/apache/seatunnel/commit/2c2ce1a543.0.0
[Fix][Connector-V2] Flush Doris load before schema change (#11156)https://github.com/apache/seatunnel/commit/bd5a19a6c3.0.0
[SEATUNNEL-10685] prevent timestamp_ntz from being saved as timestamp_ltz (#10724)https://github.com/apache/seatunnel/commit/872077f643.0.0
[Fix][Connector-V2] Fix Doris sink retry backoff and scheduler leak (#10772)https://github.com/apache/seatunnel/commit/0a2de139c3.0.0
[Feature][Connectors-v2] Add Doris sink redirect enhancement (#10715)https://github.com/apache/seatunnel/commit/6af9ed0353.0.0