Doris
Doris 源连接器
支持的引擎
Spark
Flink
SeaTunnel Zeta
主要功能
描述
用于 Apache Doris 的源连接器。
依赖
对于 Spark/Flink
- 你需要下载 jdbc driver jar package 并添加到目录
${SEATUNNEL_HOME}/plugins/.
对于 SeaTunnel Zeta
- 你需要下载 jdbc driver jar package 并添加到目录
${SEATUNNEL_HOME}/lib/.
支持的数据源信息
| 数据源 | 支持版本 | 驱动 | Url | Maven |
|---|---|---|---|---|
| Doris | 仅支持Doris2.0及以上版本. | - | - | - |
数据类型映射
| Doris 数据类型 | SeaTunnel 数据类型 |
|---|---|
| INT | INT |
| TINYINT | TINYINT |
| SMALLINT | SMALLINT |
| BIGINT | BIGINT |
| LARGEINT | STRING |
| BOOLEAN | BOOLEAN |
| DECIMAL | DECIMAL((Get the designated column's specified column size)+1, (Gets the designated column's number of digits to right of the decimal point.))) |
| FLOAT | FLOAT |
| DOUBLE | DOUBLE |
| CHAR VARCHAR STRING TEXT | STRING |
| JSON | STRING |
| VARIANT | STRING |
| DATE | DATE |
| DATETIME DATETIME(p) | TIMESTAMP |
| ARRAY | ARRAY |
源选项
基础配置:
| 名称 | 类型 | 是否必须 | 默认值 | 描述 |
|---|---|---|---|---|
| fenodes | string | yes | - | FE 地址, 格式:"fe_host:fe_http_port" |
| username | string | yes | - | 用户名 |
| password | string | yes | - | 密码 |
| doris.request.retries | int | no | 3 | 请求Doris FE的重试次数 |
| doris.request.read.timeout.ms | int | no | 30000 | 请求 Doris BE 的 socket 读取超时时间。 |
| doris.request.connect.timeout.ms | int | no | 30000 | 请求 Doris FE 或 BE 的连接超时时间。 |
| query-port | int | no | 9030 | Doris 查询端口。 |
| doris.request.query.timeout.s | int | no | 3600 | Doris扫描数据的超时时间,单位秒 |
| doris.request.tablet.size | int | no | Integer.MAX_VALUE | 每个 SeaTunnel split 包含的 Doris tablet 数量,最小值为 1。 |
| doris.deserialize.arrow.async | boolean | no | false | 是否异步反序列化 Arrow 数据。 |
| doris.request.retriesdoris.deserialize.queue.size | int | no | 64 | 异步反序列化 Arrow 数据时使用的队列大小。 |
| table_list | Array | no | - | 要读取的 Doris 表清单。 |
doris.request.retriesdoris.deserialize.queue.size 是当前运行时实际使用的配置名。调整异步 Arrow 反序列化队列大小时,请按这个完整名称配置。
表清单配置:
| 名称 | 类型 | 是否必须 | 默认值 | 描述 |
|---|---|---|---|---|
| database | string | yes | - | 数据库 |
| table | string | yes | - | 表名 |
| doris.read.field | string | no | - | 选择要读取的Doris表字段 |
| doris.filter.query | string | no | - | 数据过滤. 格式:"字段 = 值", 例如:doris.filter.query = "F_ID > 2" |
| doris.request.tablet.size | int | no | Integer.MAX_VALUE | 当前表每个 SeaTunnel split 包含的 Doris tablet 数量,最小值为 1。 |
| doris.batch.size | int | no | 1024 | 每次能够从BE中读取到的最大行数 |
| doris.exec.mem.limit | long | no | 2147483648 | 单个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
| Change | Commit | Version |
|---|---|---|
| [Fix][CDC][Zeta] Restore runtime schema from checkpoint after failover (#11503) | https://github.com/apache/seatunnel/commit/ec1b1b8b5 | 3.0.0 |
| [Improve][Connector-V2][Doris] Support partition cleanup for DROP_DATA (#11917) | https://github.com/apache/seatunnel/commit/2714e6e25 | 3.0.0 |
| [Improve][Connector-V2] Remove redundant imperative validation in Doris connector (#11858) | https://github.com/apache/seatunnel/commit/d8186ba2d | 3.0.0 |
| [Feature][shade]Refactor the seatunnel-shade module. (#9993) | https://github.com/apache/seatunnel/commit/4ba289595 | 3.0.0 |
| [Feature][Connector-V2][CDC] Support comment-related schema change events (#11025) | https://github.com/apache/seatunnel/commit/ba55ef965 | 3.0.0 |
| [Fix][Connector-V2] Cap decimal scale to what Doris 1.x accepts (#11690) | https://github.com/apache/seatunnel/commit/1cbff54fe | 3.0.0 |
| [Feature][Connector-V2] Support timer flush for Doris sink (#11506) | https://github.com/apache/seatunnel/commit/3ba5470a3 | 3.0.0 |
| [Fix][Connector-V2] Fix Doris stream load waitForContinue timeout when FE redirect is slow (#11330) | https://github.com/apache/seatunnel/commit/e6c097391 | 3.0.0 |
| [Feature][Doris] Support VARIANT source and sink mapping (#10855) | https://github.com/apache/seatunnel/commit/2c2ce1a54 | 3.0.0 |
| [Fix][Connector-V2] Flush Doris load before schema change (#11156) | https://github.com/apache/seatunnel/commit/bd5a19a6c | 3.0.0 |
| [SEATUNNEL-10685] prevent timestamp_ntz from being saved as timestamp_ltz (#10724) | https://github.com/apache/seatunnel/commit/872077f64 | 3.0.0 |
| [Fix][Connector-V2] Fix Doris sink retry backoff and scheduler leak (#10772) | https://github.com/apache/seatunnel/commit/0a2de139c | 3.0.0 |
| [Feature][Connectors-v2] Add Doris sink redirect enhancement (#10715) | https://github.com/apache/seatunnel/commit/6af9ed035 | 3.0.0 |