
https://github.com/apache/
点击蓝字
关注我们
如果外部系统不能直连数据库,只有调用HTTP服务才能接收数据,SeaTunnel就能搞定!这篇教程带你从零搭建MySQL到WEB系统的HTTP API进行数据同步。
环境准备:一键安装
config/plugin_config文件,加入两个核心连接器:connector-http-baseconnector-cdc-mysql
sh bin/install-plugin.sh
mysql-connector-java-8.0.28.jar(或其他8.x版本)复制到SeaTunnel的lib/目录。配置MySQL:开启“数据源”
实时同步的源头是MySQL的Binlog,需要先打开它。
找到MySQL的配置文件(my.cnf或my.ini),加入下面几行,然后重启MySQL服务:
[mysqld]server-id = 1 # 集群内唯一编号log-bin = mysql-bin # 开启二进制日志binlog-format = ROW # 必须设为ROW模式
架起接收端:一个简单的HTTP服务
我们得有个地方接收数据,这里用Go快速写一个。创建一个server.go文件:
package mainimport ( "fmt" "io" "log" "net/http")// 处理所有HTTP请求,打印接收到的数据func handler(w http.ResponseWriter, r *http.Request) { defer r.Body.Close() // 记得关闭请求体 body, err := io.ReadAll(r.Body) if err != nil { msg := fmt.Sprintf("读取请求体失败: %v", err) http.Error(w, msg, http.StatusInternalServerError) fmt.Println(msg) return } if len(body) > 0 { fmt.Printf("请求内容: %s\n", string(body)) fmt.Printf("数据大小: %d 字节\n", len(body)) } else { fmt.Println("请求体为空") } // 给客户端一个成功回复 w.WriteHeader(http.StatusOK) fmt.Fprintf(w, "数据接收成功! 大小: %d 字节。", len(body))}func main() { http.HandleFunc("/", handler) // 所有请求都交给handler处理 fmt.Println("HTTP服务启动,端口 9090...") fmt.Println("你可以用curl测试: curl -X POST -d '{\"test\":123}' http://localhost:9090/") if err := http.ListenAndServe(":9090", nil); err != nil { log.Fatal("启动服务失败: ", err) }}启动它:
go run server.go
看到“HTTP服务启动...”的提示,就说明服务已经在本地9090端口待命了。
准备测试数据
在你的MySQL数据库里,建一张测试表,并插入两条初始数据:
CREATE TABLE `post` ( `id` int(11) NOT NULL, `content` varchar(50) DEFAULT NULL, `author` varchar(50) DEFAULT NULL, PRIMARY KEY (`id`)) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;INSERT INTO `post` VALUES (1, 'Mysql入门', '张三');INSERT INTO `post` VALUES (2, 'Http解析', '李四');
核心:编写SeaTunnel任务配置
在job/目录下,新建mysqlcdc_http.conf,内容就是我们的同步脚本:
env { parallelism = 1 job.mode = "STREAMING" # 流式任务模式 checkpoint.interval = 10000 # 10秒做一次检查点}source { MySQL-CDC { username = "root" password = "root" table-names = ["your_database.post"] # 改成你的数据库和表名 url = "jdbc:mysql://你的MySQL地址:3306/your_database" # 修改为你的连接信息 schema-changes.enabled = true # 允许表结构变更 }}sink { Http { url = "http://你的服务地址:9090/" # 改成你的HTTP服务地址 headers { Accept="application/json", Content-Type="application/json;charset=utf-8" } }}注意:把里面的数据库连接信息和HTTP地址替换成你自己的。
启动同步,见证时刻
在SeaTunnel根目录下,执行命令启动任务:
bin/seatunnel.sh --config job/mysqlcdc_http.conf -m local
实时效果验证
一旦任务启动成功,你会立刻在刚才的Go服务窗口看到日志:
E:\>go run server.goStarting HTTP server on port 9090...Send POST/PUT requests with body data to test.Example: curl -X POST -d '{"message":"test"}' http://localhost:9090/Request Body Content: {"id":1,"content":"Mysql入门","author":"张三"}Body Length: 50 bytesRequest Body Content: {"id":2,"content":"Http解析","author":"李四"}Body Length: 49 bytes这说明启动时,表中已有的两条数据已经被立刻捕获并推送到HTTP服务了!
为了验证实时性,现在回到MySQL,插入一条新数据:
INSERT INTO `post` VALUES (3, 'SeaTunnel实时同步测试', '王五');
数据就会打印出来
Request Body Content: {"id":3,"content":"SeaTunnel实时同步测试","author":"王五"}Body Length: 66 bytes综上,这套从 MySQL CDC 到 SeaTunnel API Sink 再到 自定义HTTP服务 的实时数据链路。以后往别的系统实时发送数据,不用直连数据库,通过HTTP服务就搞定,在跨业务,跨部门之间的数据通讯多了一个选择。
来源 | 数仓生态圈
Apache SeaTunnel是一个云原生的多模态、高性能海量数据集成工具。北京时间 2023 年 6 月1 日,全球最大的开源软件基金会ApacheSoftware Foundation正式宣布SeaTunnel毕业成为Apache顶级项目。目前,SeaTunnel在GitHub上Star数量已达9k+。SeaTunnel支持在云数据库、本地数据源、SaaS、大模型等170多种数据源之间进行数据实时和批量同步,支持CDC、DDL变更、整库同步等功能,更是可以和大模型打通,让大模型链接企业内部的数据。
同步Demo
新手入门

最佳实践

测试报告

源码解析



