MySQL 数据实时推给三方业务?试试用 SeaTunnel 搞定 HTTP 接入

如果外部系统不能直连数据库,只有调用HTTP服务才能接收数据,SeaTunnel就能搞定!这篇教程带你从零搭建My...
17836784274809e6752b03a3968dd

https://github.com/apache/SeaTunnel

点击蓝字



关注我们

如果外部系统不能直连数据库,只有调用HTTP服务才能接收数据,SeaTunnel就能搞定!这篇教程带你从零搭建MySQL到WEB系统的HTTP API进行数据同步。

环境准备:一键安装

安装必需插件
编辑config/plugin_config文件,加入两个核心连接器:

connector-http-baseconnector-cdc-mysql
然后执行安装脚本:


sh bin/install-plugin.sh
提示:如果你网络不好,也可以去Maven中央仓库手动下载对应版本的JAR包,直接放进connectors/目录。
准备MySQL驱动
mysql-connector-java-8.0.28.jar(或其他8.x版本)复制到SeaTunnel的lib/目录。

配置MySQL:开启“数据源”

实时同步的源头是MySQL的Binlog,需要先打开它。

找到MySQL的配置文件(my.cnfmy.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

Apache SeaTunnel是一个云原生的多模态、高性能海量数据集成工具。北京时间 2023 年 6 月1 日,全球最大的开源软件基金会ApacheSoftware Foundation正式宣布SeaTunnel毕业成为Apache顶级项目。目前,SeaTunnel在GitHub上Star数量已达9k+。SeaTunnel支持在云数据库、本地数据源、SaaS、大模型等170多种数据源之间进行数据实时和批量同步,支持CDC、DDL变更、整库同步等功能,更是可以和大模型打通,让大模型链接企业内部的数据。




同步Demo

MySQL→Doris | MySQLCDC | MySQL→Hive | HTTP → Doris | HTTP → MySQL | MySQL→StarRocks|MySQL→Elasticsearch |Kafka→ClickHouse

新手入门

SeaTunnel 让数据集成变得 So easy!/ 3 分钟入门指南
0 到 1 快速入门 /初探/深入理解
分布式集群部署 | CDC数据同步管道 | Oracle-CDC
图片

最佳实践

中控技术天翼云多点OPPO | 清风马蜂窝孩子王哔哩哔哩唯品会众安保险兆原数通 | 亚信科技|映客|翼康济世|信也科技|华润置地|Shopee|京东科技|58同城|互联网银行|JPMorgan
图片

测试报告

SeaTunnel VS GLUE | VS Airbyte | VS DataX|SeaTunnel 与 DataX 、Sqoop、Flume、Flink CDC 对比
图片

源码解析

Zeta引擎源码解析(一) |(二) |(三)| API 源码解析 |2.1.1源码解析|封装 Flink 连接数据库解析





仓库地址:
https://github.com/apache/seatunnel
网址:
https://seatunnel.apache.org/
Apache SeaTunnel 下载地址:
https://seatunnel.apache.org/download
衷心欢迎更多人加入!
我们相信,在Community Over Code(社区大于代码)、「Open and Cooperation」(开放协作)、「Meritocracy」(精英管理)、以及「多样性与共识决策」The Apache Way 的指引下,我们将迎来更加多元化和包容的社区生态,共建开源精神带来的技术进步!
我们诚邀各位有志于让本土开源立足全球的伙伴加入 SeaTunnel 贡献者大家庭,一起共建开源!
提交问题和建议:
https://github.com/apache/seatunnel/issues
贡献代码:
https://github.com/apache/seatunnel/pulls
订阅社区开发邮件列表 :
dev-subscribe@seatunnel.apache.org
开发邮件列表:
dev@seatunnel.apache.org
加入 Slack:
https://join.slack.com/t/apacheseatunnel/shared_invite/zt-3uouszk3m-PtLLNyZsJVqE5Gb6gn24mA
关注 X.com:
https://x.com/ASFSeaTunnel


1783678429715dfb380e33dab0cdb
17836784304405d4ce5255f601542
17836784316171d1e5050dd435c27