☰
SeaTunnel My Hours Source Connector 实战指南:登录鉴权、配置参数与 JSON 数据抽取原理
2026/9/29 3:33:43 网站建设 项目流程
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/gh_mirrors/sea/seatunnel
点击查看免费下载

My Hours 是一款面向团队工时与项目管理的时间追踪 SaaS 服务,本指南围绕 SeaTunnel 社区版中新增的MyHoursSource 连接器展开,讲解如何通过其开放 REST API 将 My Hours 中的项目、预算与工时数据批量抽取进 SeaTunnel 数据管道。读完本文,你将掌握 My Hours 连接器的完整配置语法、登录令牌获取机制、format/content_field/json_field三类 JSON 抽取模式的实战用法,以及底层基于 HTTP 连接器框架的实现原理与排错方法。

连接器概览:My Hours 在 SeaTunnel 中的定位

MyHours是 SeaTunnel Connector V2(HTTP 连接器家族)中的一个 Source 插件,用于读取 My Hours 服务中的数据。它与 GitHub、Jira 等同类连接器一样,本质上是带特定鉴权逻辑的 HTTP 源:先通过账号密码换取访问令牌,再携带令牌请求业务 API,最后把返回的 JSON 响应反序列化为 SeaTunnel 行数据。

从仓库结构看,该连接器位于:

  • 插件实现:seatunnel-connectors-v2/connector-http/connector-http-myhours
  • 共享 HTTP 基础能力:seatunnel-connectors-v2/connector-http/connector-http-base

其插件标识符(factoryIdentifier)为MyHours,对应seatunnel.source.MyHours = connector-http-myhours的插件映射(见 plugin-mapping.properties),模块坐标由 connector-http-myhours/pom.xml 声明为connector-http-myhours(归属于org.apache.seatunnel:connector-http父模块,依赖connector-http-base)。

依赖获取

使用 My Hours 连接器需要引入如下依赖,可通过install-plugin.sh脚本离线安装,或从 Maven 中央仓库获取org.apache.seatunnel/seatunnel-connectors-v2系列构件:

DatasourceSupported VersionsDependency
My Hoursuniversalseatunnel-connectors-v2(含 connector-http-myhours 子模块)

说明:安装插件后,在任务配置的source块中以MyHours { ... }形式声明即可被 SeaTunnel 引擎识别加载。

支持引擎与功能特性

该连接器支持以下三种运行引擎:

Spark、Flink、SeaTunnel Zeta

当前版本的功能特性如下(对照 Connector V2 特性说明):

特性是否支持
batch 批式读取✅ 支持
stream 流式读取❌ 不支持
exactly-once 精确一次❌ 不支持
column projection 列投影❌ 不支持
parallelism 并行度❌ 不支持
support user-defined split 用户自定义分片❌ 不支持

从源码结构看,MyHoursSource继承自 HttpSource,而HttpSource继承自AbstractSingleSplitSource(单分片抽象源),这正是"不支持并行度、不支持自定义分片"特性的实现基础:整个 Source 以单分片方式读取,适合数据量可控的 API 拉取场景。

Source 参数总览

连接器的完整参数如下表(继承自 HTTP 基类的通用参数含义与 Source Common Options 一致):

参数名类型必填默认值说明
urlString是-HTTP 请求地址(My Hours 业务 API,如https://api2.myhours.com/api/Projects/getAll)
emailString是-My Hours 登录邮箱地址
passwordString是-My Hours 登录密码
schemaConfig否-HTTP 返回数据与 SeaTunnel 表结构之间的映射
schema.fieldsConfig否-上游数据的字段定义
json_fieldConfig否-配合 schema 使用,按 JsonPath 逐字段抽取上游数据
content_jsonString否-按 JsonPath 抽取整段 JSON 数据(如只需book段可配content_field = "$.store.book.*")
formatString否json上游数据格式,目前仅支持json、text
methodString否getHTTP 请求方法,仅支持 GET、POST
headersMap否-HTTP 请求头
paramsMap否-HTTP 请求参数
bodyString否-HTTP 请求体
poll_interval_millisInt否-流式模式下轮询 HTTP API 的间隔(毫秒)
retryInt否-请求返回IOException时的最大重试次数
retry_backoff_multiplier_msInt否100请求失败时的退避时间(毫秒)倍数
retry_backoff_max_msInt否10000请求失败时的最大退避时间(毫秒)
enable_multi_linesBoolean否false是否按行切分 HTTP 响应文本
common-options-否-Source 插件通用参数,详见 Source Common Options

关于content_json的命名说明:参数表中写作content_json,但当前仓库源码(HttpConfig.java 中CONTENT_FIELD定义)与文档示例实际使用的键名均为content_field。配置时请以content_field为准,例如content_field = "$.store.book.*"。

除上表外,从源码还可确认连接器继承了两个未在上表列出的超时参数,可供精细化调优:

  • connect_timeout_ms:连接超时,默认 12 秒(DEFAULT_CONNECT_TIMEOUT_MS = 6000 * 2)
  • socket_timeout_ms:Socket 超时,默认 60 秒(DEFAULT_SOCKET_TIMEOUT_MS = 6000 * 10)

两者在 HttpConfig.java 与 HttpParameter.java 中定义并被HttpClientProvider的RequestConfig采用。

必填参数与登录鉴权机制:email / password / url

MyHours连接器有三个必填参数:url、email、password。在 MyHoursSource.java 的构造函数中,会通过CheckConfigUtil.checkAllExists对三者做存在性校验,任一缺失即抛出配置校验异常:

CheckResult result = CheckConfigUtil.checkAllExists( pluginConfig, MyHoursSourceConfig.URL.key(), MyHoursSourceConfig.EMAIL.key(), MyHoursSourceConfig.PASSWORD.key()); if (!result.isSuccess()) { throw new MyHoursConnectorException( SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED, ...); }

登录换取 accessToken 的实现细节

与普通 HTTP 源不同,My Hours API 需要携带 Bearer 令牌访问。连接器在初始化阶段会自动完成一次"登录"请求(见MyHoursSource.getAccessToken()):

  1. 拼接登录请求:由 MyHoursSourceParameter.buildWithLoginConfig() 完成,固定使用:

    • 登录地址:https://api2.myhours.com/api/tokens/login(定义于MyHoursSourceConfig.AUTHORIZATION_URL)
    • 请求方法:POST
    • 请求体(JSON):{"grantType": "password", "email": "<邮箱>", "password": "<密码>", "clientId": "api"}
  2. 发送并解析响应:通过HttpClientProvider.doPost(url, body)发送;若响应状态码为200,则将响应体解析为 Map 并取出accessToken字段。

  3. 注入鉴权头:登录成功后,buildWithConfig()会把令牌拼成Authorization: Bearer <accessToken>写入请求头(MyHoursSourceParameter.java 中的ACCESS_TOKEN_PREFIX = "Bearer"),后续所有业务 API 请求都携带该令牌。

  4. 失败即终止:任何一步失败(非 200 响应、响应体为空、解析异常)都会抛出GET_MYHOURS_TOKEN_FAILE(错误码MYHOURS-01,"Get myhours token failed"),详见 MyHoursConnectorErrorCode.java。

因此,即使你配置的url是业务 API(而非登录接口),连接器也会先用email/password自动完成鉴权,这是理解整个连接器行为的关键。

配置一个 My Hours 数据同步作业

以下是一个完整的批式同步作业示例:从 My Hours 读取项目列表,并把结果打印到控制台。该示例完整复刻了官方文档示例,并补充了format = "json"显式声明(原因见下文"format 参数"一节的源码说明):

env { parallelism = 1 job.mode = "BATCH" } source { MyHours { url = "https://api2.myhours.com/api/Projects/getAll" email = "seatunnel@test.com" password = "seatunnel" format = "json" schema { fields { name = string archived = boolean dateArchived = string dateCreated = string clientName = string budgetAlertPercent = string budgetType = int totalTimeLogged = double budgetValue = double totalAmount = double totalExpense = double laborCost = double totalCost = double billableTimeLogged = double totalBillableAmount = double billable = boolean roundType = int roundInterval = int budgetSpentPercentage = double budgetTarget = int budgetPeriodType = string budgetSpent = string id = string } } } } # 将读取到的数据打印到控制台 sink { Console { parallelism = 1 } }

要点说明:

  • env.job.mode = "BATCH"决定了 Source 的有界性(BOUNDED)。从 HttpSource.getBoundedness() 的实现看:批模式返回Boundedness.BOUNDED,读完后通过context.signalNoMoreElement()结束;流式模式(job.mode = "STREAMING")返回Boundedness.UNBOUNDED,配合poll_interval_millis持续轮询。
  • schema.fields中声明的字段必须与 API 实际返回的 JSON 字段一一对应,类型按 SeaTunnel 类型系统声明(string / int / double / boolean)。
  • 若想让下游 Transform / Sink 引用本 Source 产出的数据,可补充result_table_name,并在下游声明source_table_name,规则详见 Source Common Options。

format 参数:json 与 text 两种解析模式

format控制上游响应数据的解析方式,仅支持json和text两种取值(对应源码 HttpConfig.ResponseFormat 枚举)。

format = json(需配合 schema)

当指定format = "json"时,必须同时配置schema用于反序列化。例如上游数据为:

{ "code": 200, "data": "get success", "success": true }

则 schema 应配置为:

schema { fields { code = int data = string success = boolean } }

连接器将生成如下数据行:

codedatasuccess
200get successtrue

从源码看,format = "json"分支会构造JsonDeserializationSchema(见 HttpSource.buildSchemaWithConfig()),并额外支持json_field与content_field两种预处理(详见后文)。

format = text(原样透传)

当指定format = "text"时,连接器不对上游数据做任何解析,整段响应文本作为一条记录输出。例如上游数据仍是:

{ "code": 200, "data": "get success", "success": true }

连接器将生成:

content
{"code": 200, "data": "get success", "success": true}

此时 Source 输出的表结构是单个content字符串列(由 HttpSource.java 中"未配置 schema"的默认分支创建,并采用SimpleTextDeserializationSchema)。

源码级注意事项:参数表标注format默认值为json,但当前仓库中HttpConfig.FORMAT的实际默认值为TEXT,且带 schema 场景的解析逻辑switch仅实现了JSON分支,其余取值会抛出Unsupported data format异常。因此只要声明了schema,就应显式配置format = "json"(仓库 e2e 测试配置如 http_jsonpath_to_assert.conf、http_contentjson_to_assert.conf 也均显式声明了format = "json")。

content_field:按 JsonPath 抽取局部 JSON 数据

当 My Hours API 返回的 JSON 体积较大、而你只需要其中某一段时,可用content_field(参数表写作content_json)按 JsonPath 抽取。例如只需要store下的book数组段:

{ "store": { "book": [ { "category": "reference", "author": "Nigel Rees", "title": "Sayings of the Century", "price": 8.95 }, { "category": "fiction", "author": "Evelyn Waugh", "title": "Sword of Honour", "price": 12.99 } ], "bicycle": { "color": "red", "price": 19.95 } }, "expensive": 10 }

配置content_field = "$.store.book.*"后,返回结果被裁剪为:

[ { "category": "reference", "author": "Nigel Rees", "title": "Sayings of the Century", "price": 8.95 }, { "category": "fiction", "author": "Evelyn Waugh", "title": "Sword of Honour", "price": 12.99 } ]

此时配合更精简的 schema 即可得到期望结果:

source { Http { url = "http://mockserver:1080/contentjson/mock" method = "GET" format = "json" content_field = "$.store.book.*" schema = { fields { category = string author = string title = string price = string } } } }

实现层面,HttpSourceReader.collect() 会先调用getPartOfJson()(内部通过JsonPath.using(jsonConfiguration).parse(data).read(JsonPath.compile(contentJson))求值),把裁剪后的 JSON 数组再交给JsonDeserializationSchema逐条反序列化。

仓库中配套的 e2e 测试佐证如下:

  • Mock 数据:mockserver-config.json
  • 任务配置:http_contentjson_to_assert.conf

json_field:多字段 JsonPath 映射

json_field用于"按字段逐个指定 JsonPath 抽取路径",它必须与schema配合使用。仍以上面的store.book数据为例,可以通过如下配置只抽取 book 段并按字段重组:

source { Http { url = "http://mockserver:1080/jsonpath/mock" method = "GET" format = "json" json_field = { category = "$.store.book[*].category" author = "$.store.book[*].author" title = "$.store.book[*].title" price = "$.store.book[*].price" } schema = { fields { category = string author = string title = string price = string } } } }

其底层原理(见 HttpSourceReader.java):

  • 使用 Jayway JsonPath,开启SUPPRESS_EXCEPTIONS、ALWAYS_RETURN_LIST、DEFAULT_PATH_LEAF_TO_NULL三个选项,保证路径不存在时不抛异常、结果恒为列表;
  • 各字段的 JsonPath 分别求值后按行"翻转"(dataFlip)重组为一行行Map<字段名, 值>;
  • 若某条 JsonPath 求值出的记录数与首个字段不一致,会抛出FIELD_DATA_IS_INCONSISTENT异常,提示解析记录数不一致——这提醒我们各 JsonPath 必须指向等长的数组。

仓库中配套的 e2e 测试配置为 http_jsonpath_to_assert.conf,Mock 数据同样位于 mockserver-config.json。

HTTP 请求细节参数与重试机制

method / headers / params / body

  • method:仅支持GET、POST(与当前仓库 HttpRequestMethod 枚举完全一致),默认GET;
  • headers/params:以 Map 形式配置,HttpParameter.buildWithConfig()会将其逐项解析为请求头与 URL 参数;对于 My Hours,Authorization: Bearer <token>由连接器自动注入,无需手工配置;
  • body:POST 请求体字符串;HttpClientProvider在携带 body 时会自动设置Content-Type: application/json。

retry 系列参数

  • retry:请求抛出IOException时的最大重试次数。源码(HttpClientProvider.buildRetryer())中仅对IOException触发重试(retryIfException),并采用stopAfterAttempt(retry)停止策略;
  • retry_backoff_multiplier_ms(默认 100):退避等待采用Fibonacci 策略(fibonacciWait(multiplier, max, MILLISECONDS)),即第 N 次重试前等待的时间按斐波那契数列放大;
  • retry_backoff_max_ms(默认 10000):退避等待时间的上限。

poll_interval_millis 与流式模式

poll_interval_millis用于流式模式下的轮询间隔。在 HttpSourceReader.internalPollNext() 中,若任务为无界(UNBOUNDED)且该值大于 0,则每次拉取后Thread.sleep(pollIntervalMillis)再进入下一轮。注意:本文档标注 My Hours 连接器不支持 stream 特性,此参数更多是为 HTTP 基类能力预留。

enable_multi_lines

enable_multi_lines = true时,连接器会把 HTTP 响应文本按行切分(BufferedReader.readLine()),每一行作为独立记录进入下游解析流程(见 HttpSourceReader.pollAndCollectData()),适合响应体为多行 JSON 序列的场景。

分页扩展(从基类继承)

从源码结构看,HTTP 基类还实现了pageing分页能力(HttpSource.buildPagingWithConfig()):通过total_page_size、batch_size(默认 100)、page_field(默认page)自动递增页码请求,直到读取条数小于batch_size或达到总页数。My Hours 连接器继承了这一能力,但本文档参数表未单独列出,属可选高级用法,可按需探索。

测试与验证

仓库中为 My Hours 连接器提供了单元测试 MyHoursFactoryTest.java,验证MyHoursSourceFactory.optionRule()能正常构建参数规则(即必填项url、email、password与可选项的声明完整无误)。

而content_field与json_field两类抽取路径的正确性,由 HTTP 连接器家族的 e2e 测试覆盖:

  • http_contentjson_to_assert.conf:验证content_field抽取后各字段非空;
  • http_jsonpath_to_assert.conf:验证json_field多字段 JsonPath 映射后各字段非空;
  • Mock 服务数据统一由 mockserver-config.json 提供(基于 MockServer 的 request matcher 机制)。

常见问题排查

现象可能原因与处理
作业启动报MYHOURS-01(Get myhours token failed)登录换取 accessToken 失败。请检查email、password是否正确、My Hours 登录接口https://api2.myhours.com/api/tokens/login是否可达、账号是否有 API 访问权限;错误信息中会附带 HTTP 状态码与响应内容
配置校验失败(CONFIG_VALIDATION_FAILED)缺少url/email/password三者之一,构造函数中的checkAllExists会给出缺失项提示
抛Unsupported data format异常声明了schema但未显式设置format = "json"(当前仓库format默认值为text,而 schema 场景仅支持 json),请显式声明format = "json"
抛FIELD_DATA_IS_INCONSISTENT异常json_field中各 JsonPath 求值出的记录数不一致(如某字段路径指向长度不同的数组),请核对各 JsonPath 表达式
请求频繁失败可配置retry及退避参数、适当调大connect_timeout_ms/socket_timeout_ms(默认分别为 12s / 60s)

Changelog 与演进

当前仓库的 My Hours 连接器演进记录如下:

  • next version:新增 My Hours Source Connector;HTTP 连接器引入 json-path 解析能力(对应上游 PR 3510 的[Feature][Connector-V2][HTTP] Use json-path parsing改动)。

json-path 解析能力的引入正是content_field与json_field两个参数得以实现的基础,也使得 My Hours 连接器在面对嵌套 JSON 响应时无需编写复杂预处理即可按路径精准抽取数据。

小结

My Hours 连接器是 SeaTunnel HTTP 连接器家族中"带鉴权逻辑"的代表性 Source:通过email/password自动完成令牌换取、以 Bearer 头访问业务 API,再通过format/content_field/json_field的组合实现对 JSON 响应的灵活抽取。理解其参数表、登录流程与 HttpSourceReader 的数据处理链路,即可在 Spark、Flink、SeaTunnel Zeta 三种引擎上快速搭建 My Hours 数据同步管道,并将同类 HTTP API 的接入经验复用到其他连接器上。

  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/gh_mirrors/sea/seatunnel
点击查看免费下载

相关推荐

上一篇:Mediago 前端异步优化:基于依赖的并行化(Dependency-Based Parallelization)实战指南
下一篇:使用 Claude Code Game Studios 的 Prototyper Agent 快速验证游戏机制:从假设到 PROCEED/PIVOT/KILL 决策的完整工作流

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询