连接器类型概览
SeaTunnel 的连接器生态非常丰富,基本覆盖了主流的数据存储和计算组件。无论是关系型数据库、NoSQL 还是大数据组件,都能找到对应的 Source 和 Sink 实现。
| Source | Sink |
|---|---|
| Clickhouse | Clickhouse |
| Elasticsearch | Elasticsearch |
| FakeSource | FakeSource |
| Ftp | Ftp |
| Github/Gitlab | Github/Gitlab |
| Greenplum | Greenplum |
| Hdfs file | Hdfs file |
| Hive | Hive |
| Http | Http |
| Hudi/Iceberg | Hudi/Iceberg |
| JDBC | JDBC |
| Kudu | Kudu |
| MongoDB | MongoDB |
| Mysql / MySQL CDC | Mysql / MySQL CDC |
| Redis | Redis |
| Kafka | Kafka |
| StarRocks | StarRocks |
| Phoenix | Phoenix |
| ... | ... |
MySQL 到 MySQL 全量同步
这是最基础也是最常用的场景。配置时主要关注 Source 端的查询语句和 Sink 端的写入策略。
核心参数包括连接 URL、驱动类名以及用户凭证。在 Source 端,query 字段决定了读取范围;Sink 端则通过 query 指定插入逻辑,支持预编译语句以提高性能。
env {
execution.parallelism = 2
job.mode = "BATCH"
}
source {
Jdbc {
url = "jdbc:mysql://127.0.0.1:3306/test"
driver = "com.mysql.cj.jdbc.Driver"
connection_check_timeout_sec = 100
user = "user"
password = "password"
query = "select * from base_region limit 4"
}
}
transform {
# 此处可配置 SQL 转换插件
}
sink {
jdbc {
url = "jdbc:mysql://127.0.0.1:3306/dw"
driver = "com.mysql.cj.jdbc.Driver"
user = "user"
password = "password"
query = "insert into base_region(id,region_name) values(?,?)"
}
}
启动作业只需指定配置文件路径:
./bin/seatunnel.sh --config ./config/mysql2mysql_batch.conf
MySQL 到 Hive 同步
将数据同步至数仓通常涉及引擎依赖问题。如果使用 Spark 或 Flink 引擎,需确保环境已集成 Hive;若使用 SeaTunnel Zeta 引擎,则需手动补充相关 Jar 包。
Zeta 引擎依赖:
- seatunnel-hadoop3-3.1.4-uber.jar
- hive-exec-2.3.9.jar
请将上述文件放入 $SEATUNNEL_HOME/lib/ 目录下。
配置示例如下,注意 metastore_uri 需指向正确的 Hive Metastore 地址。
env {
job.mode = "BATCH"
}
source {
Jdbc {
url = "jdbc:mysql:///127.0.0.1:3306/dw?allowMultiQueries=true&characterEncoding=utf-8"
driver = "com.mysql.cj.jdbc.Driver"
user = "user"
password = "password"
query = "select * from source_user"
}
}
sink {
Hive {
table_name = "ods.sink_user"
metastore_uri = "thrift://bigdata101:9083"
}
}
运行命令与之前类似:
./bin/seatunnel.sh --config ./config/mysql2hive.conf
增量同步策略
生产环境中往往需要按时间维度进行增量抽取。这里以订单明细表为例,根据 create_time 字段过滤数据。
关键点:
- 源端查询需包含时间范围判断。
- 利用 Shell 变量传递日期参数(如
${etl_dt}),便于调度系统调用。 - 注意时间格式转换,避免字符串比较错误。
env {
execution.parallelism = 2
job.mode = "BATCH"
}
source {
Jdbc {
url = "jdbc:mysql://127.0.0.1:3306/test"
driver = "com.mysql.cj.jdbc.Driver"
connection_check_timeout_sec = 100
user = "user"
password = "password"
query = "select * from t_order_detail where create_time >= REPLACE('\"${etl_dt}\"', 'T', ' ') and create_time < date_add(REPLACE('\"${etl_dt}\"', 'T', ' '),interval 1 day);"
}
}
sink {
jdbc {
url = "jdbc:mysql://127.0.0.1:3306/test"
driver = "com.mysql.cj.jdbc.Driver"
connection_check_timeout_sec = 100
user = "user"
password = "password"
query = "insert into ods_t_order_detail_di (id,order_id,sku_id,sku_name,img_url,order_price,sku_num,create_time) values(?,?,?,?,?,?,?,?)"
}
}
执行时通过 -i 参数传入日期:
./bin/seatunnel.sh --config ./config/mysql2mysql_ods_t_order_detail_di.conf -i etl_dt='2024-02-05'
实时同步 (CDC)
对于对时效性要求极高的场景,可以开启流式模式 (STREAMING)。基于 MySQL CDC 可以实现准实时的数据变更捕获。
注意事项:
- 必须设置
checkpoint.interval以保证故障恢复能力。 - Sink 端建议开启
generate_sink_sql = true,让框架自动生成插入语句,减少维护成本。
env {
execution.parallelism = 2
job.mode = "STREAMING"
checkpoint.interval = 10000
}
source {
MySQL-CDC {
username = "user"
password = "password"
table-names = ["test.source_user"]
base-url = "jdbc:mysql://127.0.0.1:3306/test"
}
}
sink {
jdbc {
url = "jdbc:mysql://127.0.0.1:3306/dw"
driver = "com.mysql.cj.jdbc.Driver"
username = "user"
password = "password"
generate_sink_sql = true
database = "dw"
table = "source_user_01"
primary_keys = ["userid"]
}
}
启动脚本:
./bin/seatunnel.sh --config ./config/mysql2mysql_rt.conf

