SeaTunnel 实战:用 Http Source + JDBC Sink 搭建 HTTP API 到关系型数据库的数据同步链路
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本篇实战指南以 Apache SeaTunnel 的 Http 到 JDBC 官方食谱 为骨架,完整讲解「HTTP API 拉取结构化数据 → 写入 PostgreSQL」这条端到端链路的搭建方法:从前置条件、插件安装、JDBC 驱动放置,到最小可运行配置、任务运行与结果验证,再到分页、嵌套 JSON 抽取、自动建表与 Upsert 等进阶能力。读完本文,你将掌握 Http Source 与 JDBC Sink 两个连接器的核心参数语义,并能基于仓库中的 e2e 测试资源(如 http_streaming_json_to_postgresql.conf)独立复现和扩展这条链路。
链路概览:一条从 HTTP 到关系型数据库的数据管道
这条链路的形态非常简单,却覆盖了 SeaTunnel 最典型的两种连接器用法:
HTTP API(JSON)──Http Source──> SeaTunnel Row ──JDBC Sink──> PostgreSQL 表- Http Source以
GET/POST请求拉取接口数据,把 JSON 响应按schema声明反序列化成结构化行(SeaTunnelRow); - JDBC Sink通过数据库厂商提供的 JDBC 驱动,把上游行写入关系型数据库,支持自动生成 SQL、自动建表与主键 Upsert。
当你想从 HTTP API 拉取结构化数据,并把结果落到关系型数据库中时,就可以使用这条链路。仓库中 connector-http-e2e 的测试工程就包含了一条几乎一模一样的真实用例:http_streaming_json_to_postgresql.conf(流式轮询 Http 接口并写入 PostgreSQL),可以作为对照参考。
前置条件
1. 先跑通第一个任务
链路依赖的基础环境与「跑第一个任务」完全一致。请先完成 跑第一个任务,确认本地能正常启动 SeaTunnel 并解析配置(该教程使用config/v2.batch.config.template验证安装、配置解析与执行引擎均正常)。
2. 安装链路所需插件
从 2.2.0-beta 开始,SeaTunnel 二进制发行包不再默认附带全部连接器依赖,需要按需安装。安装前先把config/plugin_config收敛成下面这样,只保留本链路需要的connector-http-base与connector-jdbc:
--seatunnel-connectors-- connector-http-base connector-jdbc --end--然后执行安装脚本并确认插件 JAR 已经落盘:
cd "${SEATUNNEL_HOME}" sh bin/install-plugin.sh ls connectors | rg 'connector-(http-base|jdbc)'关于插件安装的详细说明(如指定版本、SEATUNNEL_MAVEN_REPOSITORY镜像等),参见 部署 > 下载连接器插件。config/plugin_config中可用的连接器名与 JAR 的对应关系,可以在${SEATUNNEL_HOME}/connectors/plugins-mapping.properties(仓库根目录的 plugin-mapping.properties 为同源映射)中查到。
3. 放置 JDBC 驱动
SeaTunnel 不会统一内置所有 JDBC 驱动,原因在于不同厂商驱动的许可证、再分发条款不同,且驱动版本必须同时兼容目标数据库与 Java 运行时。本篇使用 Zeta 引擎(本地模式),需要把 PostgreSQL JDBC 驱动 JAR 放入${SEATUNNEL_HOME}/lib,然后确认已落盘:
ls "${SEATUNNEL_HOME}/lib" | rg 'postgresql'如果你使用 Spark / Flink 引擎,驱动要放到每个执行节点的
${SEATUNNEL_HOME}/plugins/Jdbc/lib/;Zeta 引擎放到${SEATUNNEL_HOME}/lib/后还需重启受影响的 SeaTunnel 进程,让驱动进入类路径。常见驱动文件名:MySQL 为mysql-connector-j-8.x.x.jar,PostgreSQL 为postgresql-42.x.x.jar,Oracle 为ojdbc8.jar。
4. 先确认 HTTP 返回内容
运行任务前,先用curl看一眼接口返回。这里直接使用 Http Source 文档里的示例接口:
curl http://mockserver:1080/example/http该接口的 mock 数据定义在 mockserver-config.json(匹配GET /example/http),返回 JSON 顶层可以看到c_map、c_array、c_string、c_boolean、c_int、c_bigint等一系列字段,本篇的 schema 只取其中c_string和c_int两个顶层字段。
关键判断点:如果你的真实接口把有效数据包在更深层字段里(例如{"code":200, "data":{...}}),就要先补json_field或content_field做字段抽取,否则 schema 解析会失败。这一点的细节见下文「嵌套 JSON 与字段抽取」一节。
5. 准备 PostgreSQL 目标库
本篇配置使用了generate_sink_sql = true(自动生成建表与写入 SQL),因此需要给 sink 用户授予在publicschema 自动建表的权限:
CREATE USER test WITH PASSWORD 'test'; CREATE DATABASE test OWNER test;重新连接到test库以后,再执行:
GRANT USAGE, CREATE ON SCHEMA public TO test;最小配置:逐段解析
把下面这份配置保存为config/http-to-jdbc.conf。这是官方食谱中的最小可运行版本,本文会逐段拆解每个参数的含义。
env { parallelism = 1 job.mode = "BATCH" } source { Http { plugin_output = "http_orders" url = "http://mockserver:1080/example/http" method = "GET" format = "json" schema = { fields { c_string = string c_int = int } } } } sink { Jdbc { plugin_input = "http_orders" driver = "org.postgresql.Driver" url = "jdbc:postgresql://postgresql:5432/test?loggerLevel=OFF" username = "test" password = "test" generate_sink_sql = true database = "test" table = "public.http_orders" primary_keys = ["c_string"] batch_size = 100 } }env 块:执行环境
parallelism = 1:并发度设为 1。Http Source 当前不支持用户自定义分片(特性矩阵中「支持用户自定义分片」未勾选),且本示例接口为单页返回,单并发即可满足。job.mode = "BATCH":批模式,任务拉完全部数据后正常退出。如果希望定时轮询接口(流模式),可改为STREAMING并配合poll_interval_millis,参见 http_streaming_json_to_postgresql.conf 中的用法。
source 块:Http Source 核心参数
| 参数 | 值 | 说明 |
|---|---|---|
plugin_output | http_orders | 本节点输出的数据流命名,sink 用plugin_input指向它完成上下游对接 |
url | http://mockserver:1080/example/http | 请求 URL |
method | GET | 请求方法,仅支持GET、POST |
format | json | 上游数据格式,支持json、text、binary,默认text。设置为json时必须配套声明schema |
schema.fields | c_string、c_int | 上游数据的字段声明,SeaTunnel 依据它把 JSON 响应反序列化为结构化行 |
值得强调的是plugin_output/plugin_input这对参数:它们把 source 节点的输出与 sink 节点的输入显式连接起来,是当前推荐的数据流命名方式。从源码看,HttpSourceReader在pollAndCollectData中拿到响应后,会走DeserializationCollector按声明 schema 产出SeaTunnelRow(见 HttpSourceReader.java),sink 侧按plugin_input接收这些行。
sink 块:JDBC Sink 核心参数
| 参数 | 值 | 说明 |
|---|---|---|
driver | org.postgresql.Driver | JDBC 驱动类名,对应connectors/plugin-mapping.properties中connector-jdbc对应的驱动 |
url | jdbc:postgresql://postgresql:5432/test?loggerLevel=OFF | JDBC 连接 URL |
username/password | test/test | 数据库账号与密码 |
generate_sink_sql | true | 写入模式开关:为true时由 SeaTunnel 根据上游 schema 与 RowKind 自动生成 INSERT / 原生 UPSERT / UPDATE / DELETE,可配合 SaveMode 与自动建表;为false时你必须提供query自定义 SQL |
database | test | 自动生成 SQL 模式下的目标 database,generate_sink_sql = true时必填 |
table | public.http_orders | 自动生成 SQL 模式下的目标表;有 schema 概念的数据库必须写成xxx.xxx形式 |
primary_keys | ["c_string"] | 用于生成数据库原生 UPSERT、UPDATE、DELETE 的目标键列 |
batch_size | 100 | 每个 batch 最多缓存的行数,达到后触发 flush 写入 |
generate_sink_sql的底层约束:从 JdbcSinkFactory.java 的 OptionRule 可以看到,连接器始终要求url、driver、schema_save_mode、data_save_mode(后两者有默认值可省略);当generate_sink_sql = true时强制要求配置database(第 245 行conditional(JdbcSinkOptions.GENERATE_SINK_SQL, true, JdbcSinkOptions.DATABASE)),当其为false时强制要求配置query(第 246 行)。generate_sink_sql的定义位于 JdbcSinkOptions.java,默认值为false。因此,没有显式设置generate_sink_sql = true的任务必须提供query,两种写入模式不可混用。
运行任务
把配置保存为config/http-to-jdbc.conf后,用本地模式运行:
cd "${SEATUNNEL_HOME}" ./bin/seatunnel.sh --config ./config/http-to-jdbc.conf -m local-m local表示以本地模式运行(不启动独立 SeaTunnel 引擎集群,适用于单机验证)。
验证结果
- 运行任务,确认日志中没有 HTTP 解析错误和 JDBC DDL 错误;
- 连接 PostgreSQL,查询目标表,核对行数与 API 返回结果是否一致:
SELECT COUNT(*) FROM public.http_orders; SELECT c_string, c_int FROM public.http_orders ORDER BY c_string;如果目标表里的数据和 HTTP 返回内容一致,这条链路就是通的。使用默认 mock 返回时,查询结果里应该能看到和curl输出一致的c_string、c_int值。
常见坑与排查
官方食谱总结了四类最容易踩的坑,这里结合源码给出更具体的判断依据:
- 返回体是 JSON,但 schema 中字段名或字段类型写错了。
format = json模式下schema.fields是反序列化的依据,字段名与接口返回不一致、类型不匹配都会导致解析失败。可先用curl核对字段名,并参考 Http Source 中的完整 schema 示例(覆盖map、array、row、decimal、timestamp等类型)。 - API 数据是嵌套结构,但没有配置
content_field或json_field。此时 schema 直接落在顶层字段上会取不到值。解决方案见下节。 - 源接口有分页或限流,但作业按单页接口处理。Http Source 内置
pageing分页能力(PageNumber/Cursor两种类型),见下文「分页拉取」一节。 - JDBC sink 虽然自动建表了,但你选的主键并不能真正唯一标识一条记录。
primary_keys不仅用于 Upsert,还参与自动建表时主键约束的生成;如果选错了键,重复数据写入时会触发主键冲突或静默覆盖,导致数据不符合预期。
嵌套 JSON 与字段抽取:json_field与content_field
当返回体把有效数据包在深层字段时,两种参数任选其一:
content_field:直接抽取某个 JSON 数组或对象片段。例如返回体形如{"store":{"book":[...]}},配置content_field = "$.store.book.*"后,连接器只把book数组部分交给 schema 解析,schema 只需声明category、author、title、price等字段即可。json_field:为每个 schema 字段单独指定 JSONPath。例如json_field = { category = "$.store.book[*].category" },它必须与schema一起使用。对应的抽取逻辑在 HttpSourceReader.java 的collect()方法中:先走getPartOfJson(content_field 分支)或decodeJSON+parseToMap(json_field 分支),再交给反序列化器。
另外,当 JSON 字段缺失时默认会抛错;设置json_filed_missed_return_null = true(源码中该选项定义于 HttpSourceOptions.java)可以让缺失字段返回null。
分页拉取:pageing与请求形态
Http Source 的分页能力集中在pageing配置块,支持两种分页类型:
PageNumber(默认):用页码推进。关键子参数:page_field:请求中的分页字段名,默认page,可在headers、params、body中使用${page}占位符;start_page_number:起始页码,默认1;total_page_size:总页数,0表示未知——此时连接器会在单页返回行数小于pageing.batch_size(默认 100)时停止,对应 HttpSourceReader.java 中readSize < pageInfo.getBatchSize()的终止判断;use_placeholder_replacement:true时按${page}占位符替换(支持"10${page}"→"105"这类带前后缀的形式),false时只做按 key 的整值替换。
Cursor(游标):设置page_type = "Cursor",用cursor_field指定请求中的游标字段名,cursor_response_field指定响应中取游标的 JSONPath。源码中当响应游标为空或与当前游标相同(noMoreElementFlag置真)时结束拉取。
关于分页最终发出的请求形态,官方文档给出了一条核心经验法则:先想清楚「最终发出的 HTTP 请求长什么样」。例如GET请求下params一定会拼进 URL 查询串;POST且keep_params_as_form = false时 body 作为 JSON 发送(未配置 body 会发送空 JSON 对象{});keep_params_as_form = true时 params 并入表单 body 且 SeaTunnel 自动补application/x-www-form-urlencoded头;keep_page_param_as_http_param = true时分页字段直接写入params。详细规则与 GET/POST 表单三种示例见 Http Source 文档「分页与最终请求形态排查」。
写入模式、SaveMode 与 Upsert
JDBC Sink 有两种互斥的写入模式,使用时先二选一:
| 使用场景 | 必需配置 | 行为 |
|---|---|---|
| 由 SeaTunnel 生成 SQL | generate_sink_sql = true、database,通常还要配置table | 根据上游 schema 与 RowKind 生成 INSERT、原生 UPSERT、UPDATE、DELETE;支持 SaveMode 与自动建表 |
| 用户提供 SQL | query = "INSERT ... VALUES (?, ...)" | 完全控制目标 SQL;?参数按上游字段顺序绑定;此模式不执行 SaveMode |
本篇采用第一种模式,并依靠两个 SaveMode 参数控制表结构与数据的处理策略(均有默认值,通常可省略):
schema_save_mode:默认CREATE_SCHEMA_WHEN_NOT_EXIST(表不存在时创建、存在则跳过);其他选项RECREATE_SCHEMA(删除重建)、ERROR_WHEN_SCHEMA_NOT_EXIST(表不存在报错)、IGNORE(跳过建表逻辑)。data_save_mode:默认APPEND_DATA(保留数据追加);其他选项DROP_DATA(清空数据)、CUSTOM_PROCESSING(配合custom_sql自定义预处理)、ERROR_WHEN_DATA_EXISTS(有数据即报错)。
关于 Upsert 行为的关键机制(可参见 JDBC Sink 故障排查):
- SeaTunnel只有拿到主键/唯一键信息时才会进入 upsert/update 路径。这个 key 可以来自显式配置的
primary_keys,未配置时会尝试从上游 Catalog 元数据继承主键,再尝试第一组 unique key;仍然没有时退化为普通 INSERT。 - 当存在 key 且
enable_upsert = true(默认)时,优先使用数据库方言原生的 upsert 语句,例如 PostgreSQL 生成INSERT ... ON CONFLICT (...) DO UPDATE(若所有字段都是主键则退化为DO NOTHING)。 - 本篇把
c_string设为主键:接口的c_string是随机字符串(mock 数据中形如"WArEB"),主键能够唯一标识一条记录,适合演示 Upsert;若任务没有重复 key 数据,可以把enable_upsert设为false以加快导入。
提升吞吐的进阶选项
batch_size越大,每个 batch 缓存行数越多,flush 越少,吞吐越高,但内存占用与故障重试量也随之增加;- 非 XA 的 MySQL 批量任务,可在 JDBC URL 中加入
rewriteBatchedStatements=true提升吞吐; - PostgreSQL 大批量导入可尝试
use_copy_statement = true走COPY <table> FROM STDIN路径(要求驱动提供getCopyAPI(),且不支持MAP/ARRAY/ROW类型); - 精确一次(exactly-once)可配置
is_exactly_once = true+xa_data_source_class_name+max_retries = 0,但要求数据库与驱动都支持 XA 事务(PostgreSQL 需启用 prepared transaction),配置前务必确认前置条件。
小结
至此,你已经走通了「HTTP API → SeaTunnel → PostgreSQL」的完整链路:从插件安装、驱动放置到最小配置、运行验证,再到嵌套 JSON 抽取、分页与 Upsert 等进阶能力。这条链路是 SeaTunnel 众多「拉取 API 数据入库」场景的通用模板——把 Http Source 换成任意带鉴权头的接口、把 JDBC Sink 的 URL 换成 MySQL / Oracle / SQL Server 等其他数据库(驱动参考见 JDBC Sink 文档),即可快速复用到自己的业务中。
相关文档
- Http Source 连接器完整文档(全部源选项、format 详解、分页示例)
- JDBC Sink 连接器完整文档(写入模式、全部参数、SaveMode、exactly-once、故障排查、驱动参考)
- 部署与插件安装
- 跑第一个任务
- e2e 参考实现:http_streaming_json_to_postgresql.conf 与 mockserver-config.json
- 源码位置:HttpSourceReader.java、HttpSourceOptions.java、JdbcSinkFactory.java
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考