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 APIJSON──Http Source── SeaTunnel Row ──JDBC Sink── PostgreSQL 表Http Source以GET/POST请求拉取接口数据把 JSON 响应按schema声明反序列化成结构化行SeaTunnelRowJDBC 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.jarPostgreSQL 为postgresql-42.x.x.jarOracle 为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?loggerLevelOFF 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_outputhttp_orders本节点输出的数据流命名sink 用plugin_input指向它完成上下游对接urlhttp://mockserver:1080/example/http请求 URLmethodGET请求方法仅支持GET、POSTformatjson上游数据格式支持json、text、binary默认text。设置为json时必须配套声明schemaschema.fieldsc_string、c_int上游数据的字段声明SeaTunnel 依据它把 JSON 响应反序列化为结构化行值得强调的是plugin_output/plugin_input这对参数它们把 source 节点的输出与 sink 节点的输入显式连接起来是当前推荐的数据流命名方式。从源码看HttpSourceReader在pollAndCollectData中拿到响应后会走DeserializationCollector按声明 schema 产出SeaTunnelRow见 HttpSourceReader.javasink 侧按plugin_input接收这些行。sink 块JDBC Sink 核心参数参数值说明driverorg.postgresql.DriverJDBC 驱动类名对应connectors/plugin-mapping.properties中connector-jdbc对应的驱动urljdbc:postgresql://postgresql:5432/test?loggerLevelOFFJDBC 连接 URLusername/passwordtest/test数据库账号与密码generate_sink_sqltrue写入模式开关为true时由 SeaTunnel 根据上游 schema 与 RowKind 自动生成 INSERT / 原生 UPSERT / UPDATE / DELETE可配合 SaveMode 与自动建表为false时你必须提供query自定义 SQLdatabasetest自动生成 SQL 模式下的目标 databasegenerate_sink_sql true时必填tablepublic.http_orders自动生成 SQL 模式下的目标表有 schema 概念的数据库必须写成xxx.xxx形式primary_keys[c_string]用于生成数据库原生 UPSERT、UPDATE、DELETE 的目标键列batch_size100每个 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()方法中先走getPartOfJsoncontent_field 分支或decodeJSONparseToMapjson_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起始页码默认1total_page_size总页数0表示未知——此时连接器会在单页返回行数小于pageing.batch_size默认 100时停止对应 HttpSourceReader.java 中readSize pageInfo.getBatchSize()的终止判断use_placeholder_replacementtrue时按${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 与 UpsertJDBC Sink 有两种互斥的写入模式使用时先二选一使用场景必需配置行为由 SeaTunnel 生成 SQLgenerate_sink_sql true、database通常还要配置table根据上游 schema 与 RowKind 生成 INSERT、原生 UPSERT、UPDATE、DELETE支持 SaveMode 与自动建表用户提供 SQLquery 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 中加入rewriteBatchedStatementstrue提升吞吐PostgreSQL 大批量导入可尝试use_copy_statement true走COPY table FROM STDIN路径要求驱动提供getCopyAPI()且不支持MAP/ARRAY/ROW类型精确一次exactly-once可配置is_exactly_once truexa_data_source_class_namemax_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),仅供参考