- 数据工程
- 数据集成
- ETL
- 后端
- 大数据
【免费下载链接】airbyte
Open-source data movement for ELT pipelines and AI agents — from APIs, databases & files to warehouses, lakes, and AI applications. Both self-hosted and Cloud.
本指南以source-castor-edc(Castor EDC 临床研究数据源连接器)为核心,系统讲解 Airbyte 声明式连接器(Declarative / Low-Code CDK)的完整形态:从 Connector Builder 生成的manifest.yaml如何定义认证、数据流、分页与增量同步,到四个必填配置项的实际含义,再到单元测试与验收测试的组织方式。读完本文,你将掌握如何阅读、配置、测试与本地运行一个 manifest-only 连接器,并能将同一套模式迁移到其他声明式连接器上。
一、连接器概览:一个由 Connector Builder 生成的声明式连接器
source-castor-edc 的 README 开宗明义地指出:这是一个用 Connector Builder 构建的声明式连接器(Declarative Connector)。与传统的 Python/Java 手写连接器不同,声明式连接器不需要编写逐行的请求、解析与状态管理代码,而是通过一份结构化的 YAML 清单(manifest)声明"如何请求、如何解析、如何分页、如何增量",由 Airbyte 的 Low-Code CDK(声明式 CDK)在运行时解释执行。
在 metadata.yaml 中可以确认该连接器的技术属性:
connectorType: source、connectorSubtype: api,属于 REST API 类源连接器;tags同时标注了language:manifest-only与cdk:low-code,表明它没有语言级业务代码,全部逻辑由 manifest 承载;releaseStage: alpha、supportLevel: community,releaseDate: 2024-10-12,当前dockerImageTag: 0.0.64;connectorBuildOptions.baseImage指向airbyte/source-declarative-manifest:7.33.0,即运行时代理镜像,印证了"manifest 驱动 + 通用运行镜像"的运行模式;allowedHosts列出了三个区域主机:data.castoredc.com、uk.castoredc.com、us.castoredc.com。
连接器目录结构非常精简,全部文件如下:
| 文件/目录 | 作用 |
|---|---|
| manifest.yaml | 连接器的唯一实现:声明流、认证、配置、schema(约 2500 行) |
| metadata.yaml | 注册元数据:定义 ID、镜像名、发布阶段、允许主机等 |
| acceptance-test-config.yml | Connector Acceptance Tests(CAT)配置 |
| unit_tests/test_manifest.py | 针对 manifest 行为的单元测试 |
| unit_tests/pyproject.toml | 测试环境声明(airbyte-cdk 7.33.0、pytest ^8.0、pyyaml ^6.0) |
| icon.svg | 连接器图标 |
可见,对这种 manifest-only 连接器而言,README 只是入口,manifest.yaml 才是灵魂。下文将以 manifest.yaml 为主线逐层拆解。
二、manifest 顶层骨架:check、streams 与 spec 三大件
manifest.yaml 的顶层结构是标准声明式连接器三件套:
version: 5.12.0 type: DeclarativeSource description: "Documentation: https://uk.castoredc.com/api#/" check: type: CheckStream stream_names: - user definitions: # 所有可复用的组件定义:streams、base_requester streams: # 暴露给 Airbyte 的 16 个数据流(引用 definitions 中的流) spec: # 连接器的配置表单(connection_specification) schemas: # 每个流的内联 JSON Schematype: DeclarativeSource告诉 Low-Code CDK 这是一个声明式连接器;version: 5.12.0是 manifest 格式版本号;check阶段使用CheckStream,通过请求user流来验证连接配置是否有效——这是连接器在建立连接时执行的连通性检查;definitions是组件的"仓库",streams一节通过$ref引用其中的流,这种"定义-引用"分离的结构让多个流可以复用同一个请求器(requester)。
具体到本连接器,16 个流在 streams 列表 中逐一注册,而schemas一节为每个流提供了内联 JSON Schema(InlineSchemaLoader)。
三、认证与区域路由:base_requester 与 OAuth client_credentials
所有流共享同一个请求器base_requester,定义在 manifest.yaml 的 definitions 段:
base_requester: type: HttpRequester url_base: https://{{ 'data' if config['url_region'] == 'nl' else config['url_region'] }}.castoredc.com/api/ authenticator: type: OAuthAuthenticator client_id: "{{ config[\"client_id\"] }}" grant_type: client_credentials client_secret: "{{ config[\"client_secret\"] }}" refresh_request_body: {} token_refresh_endpoint: https://{{ 'data' if config['url_region'] == 'nl' else config['url_region'] }}.castoredc.com/oauth/token这里体现了两个值得注意的设计:
区域动态路由:
url_base使用 Jinja 模板表达式,当url_region为nl时主机为data.castoredc.com(Castor 荷兰区使用data子域),其余情况直接使用区域值(uk、us)。这与 metadata.yaml 中allowedHosts的三个主机完全对应,也呼应了 Castor 官方用户文档 中"nl同时用于 API 请求与 OAuth 认证"的说明。OAuth client_credentials 流程:认证器采用
OAuthAuthenticator,grant_type为client_credentials,通过client_id/client_secret向/oauth/token端点换取访问令牌,令牌随后自动附加到 API 请求中。refresh_request_body: {}表示刷新请求体为空。从流定义可以看到,每个流的requester都通过$ref: "#/definitions/base_requester"复用该认证逻辑。
这一行为并非笔者的主观推断,unit_tests/test_manifest.py 提供了机器可验证的证据:测试将 manifest 中的base_requester交给ModelToComponentFactory实例化为真实组件,然后断言三个区域的 API 地址与令牌端点分别为:
assert api_url == f"https://{host}/api/" assert token_url == f"https://{host}/oauth/token" assert requester.authenticator.get_grant_type() == "client_credentials"其中nl -> data.castoredc.com、uk -> uk.castoredc.com、us -> us.castoredc.com,并校验主机均在allowedHosts白名单内。第二项测试 test_region_configuration_remains_compatible 则锁定url_region的枚举值必须为["uk", "nl", "us"]且默认值为uk,防止未来改动破坏配置兼容性。
四、16 个数据流全解:三类流模式
Castor 官方用户文档 的 Streams 一节给出了 16 个流的完整清单。结合 manifest.yaml 的实现,可以进一步整理出每个流的主键、分页策略、增量能力与记录提取路径:
| 流名 | 主键 | 分页 | 增量(游标字段) | 记录提取路径(DpathExtractor) | 父流 |
|---|---|---|---|---|---|
user | id | 无 | 支持(last_login) | _embedded.user | - |
study | crf_id | DefaultPaginator(每页 50) | 支持(created_on) | _embedded.study | - |
audit_trial | uuid | DefaultPaginator(每页 1000) | 支持(date) | items | study |
country | id | 无 | 不支持 | results | - |
study_survey | id | DefaultPaginator | 不支持 | _embedded.surveys | study |
study_survey_package | id | DefaultPaginator | 不支持 | _embedded.survey_packages | study |
study_site | id | DefaultPaginator | 不支持 | _embedded.sites | study |
study_fields | id | DefaultPaginator | 不支持 | _embedded.fields | study |
study_field_dependency | id | DefaultPaginator | 不支持 | _embedded.fieldDependencies | study |
study_field_validation | id | DefaultPaginator | 不支持 | _embedded.fieldValidations | study |
study_form | id | DefaultPaginator | 不支持 | _embedded.forms | study |
study_role | uuid | DefaultPaginator | 不支持 | _embedded.roles | study |
study_fieldoption_groups | id | DefaultPaginator | 不支持 | _embedded.fieldOptionGroups | study |
study_statistics | study_id | DefaultPaginator | 不支持 | 空路径(取整个响应体) | study |
study_user | id | DefaultPaginator | 支持(last_login) | _embedded.studyUsers | study |
study_visit | id | DefaultPaginator | 不支持 | _embedded.visits | study |
从源码结构可以归纳出三类典型的流模式:
模式一:顶层列表流(无需父流)。user、study、country直接请求根级端点(如GET /api/user、GET /api/study),其中user与country甚至不配置分页器,属于一次性全量拉取;study是其余 13 个流的"数据源根"。
模式二:按 study 分区的子流(Substream)。绝大多数流属于此类,其端点路径形如study/{{ stream_partition['study_id'] }}/site、study/{{ stream_partition['study_id'] }}/form,并通过SubstreamPartitionRouter将父流study的每条记录拆分为一个分区请求。典型配置(以study_visit为例,见 manifest.yaml):
partition_router: type: SubstreamPartitionRouter parent_stream_configs: - type: ParentStreamConfig parent_key: study_id # 父流 study 中取值字段 partition_field: study_id # 注入子流路径的占位符 stream: $ref: "#/definitions/streams/study"这意味着连接器先同步study流,再为每个study_id依次请求其下的调查问卷(survey)、表单(form)、站点(site)、字段(field)、角色(role)等子资源。
模式三:统计聚合流。study_statistics的record_selector使用空field_path: [],从源码结构看,这表示不深入响应对象,而是将整个响应体作为记录输出,用于获取每个研究的记录数与按机构拆分的统计信息。
record_selector统一采用RecordSelector+DpathExtractor组合,例如user流从_embedded.user提取记录(Castor API 遵循 HATEOAS 风格,数据包裹在_embedded下),country流则从results提取,这与各端点实际的响应结构一一对应。
五、分页、限流与重试策略
除user、country外,其余流都配置了DefaultPaginator与PageIncrement分页策略。以study流为例(manifest.yaml):
paginator: type: DefaultPaginator page_token_option: type: RequestOption inject_into: request_parameter field_name: page page_size_option: type: RequestOption field_name: page_size inject_into: request_parameter pagination_strategy: type: PageIncrement page_size: 50 start_from_page: 1 inject_on_first_request: true要点解析:
- 分页参数
page与page_size以query 参数(request_parameter)方式注入请求; PageIncrement表示每页页码递增:start_from_page: 1从第 1 页开始,inject_on_first_request: true表示首次请求也携带页码参数;study流每页 50 条,而其余子流(如audit_trial、study_survey)每页 1000 条,不同流的page_size可根据接口特性差异化设置。
所有流都配置了统一的限流与重试策略(CompositeErrorHandler,例如 manifest.yaml):
error_handler: type: CompositeErrorHandler error_handlers: - type: DefaultErrorHandler max_retries: 5 response_filters: - type: HttpResponseFilter action: RATE_LIMITED http_codes: - 429 error_message: Rate limits has been hit backoff_strategies: - type: ConstantBackoffStrategy backoff_time_in_seconds: 5含义:当接口返回 429(请求过于频繁)时,连接器将命中RATE_LIMITED过滤器,以"每 5 秒固定退避"的方式最多重试 5 次。这种"识别限流码 + 固定退避"的模式是声明式连接器处理第三方 API 限流的标准做法,避免同步任务因瞬时限流而失败。
六、增量同步与时间游标:DatetimeBasedCursor
4 个流支持增量同步:user(游标last_login)、study(游标created_on)、audit_trial(游标date)、study_user(游标last_login)。它们统一使用DatetimeBasedCursor,以study流为例(manifest.yaml):
incremental_sync: type: DatetimeBasedCursor cursor_field: created_on cursor_datetime_formats: - "%Y-%m-%d %H:%M:%S" datetime_format: "%Y-%m-%d %H:%M:%S" start_datetime: type: MinMaxDatetime datetime: "{{ config[\"start_date\"] }}" datetime_format: "%Y-%m-%dT%H:%M:%SZ" end_datetime: type: MinMaxDatetime datetime: "{{ now_utc().strftime('%Y-%m-%dT%H:%M:%SZ') }}" datetime_format: "%Y-%m-%dT%H:%M:%SZ"要点解析:
- 游标字段分别为
created_on(研究创建时间)与last_login(用户最近登录时间); start_datetime取用户配置的start_date,end_datetime取当前 UTC 时间(now_utc()),构成同步窗口;audit_trial与study_user还额外通过start_time_option/end_time_option将窗口映射为请求参数date_from/date_to(见 manifest.yaml),说明这两个接口以服务端过滤配合本地游标双重实现增量。
audit_trial流的transformations部分展示了声明式字段变换的用法(manifest.yaml):
transformations: - type: AddFields fields: - path: [date] value: "{{ record[\"datetime\"][\"date\"].split(\".\")[0] }}" - type: AddFields fields: - path: [uuid] value: "{{ now_utc() }}"第一条变换从嵌套的datetime.date对象中截取日期字符串作为顶层date字段(也是增量游标);第二条为每条记录注入now_utc()生成的 UUID,充当该流的主键(primary_key: uuid)。study_role流同样使用now_utc()生成uuid主键。这体现了声明式连接器"以变换补足主键与游标"的常见技巧。
七、配置表单与使用方式:四个必填参数
连接器的配置表单定义在 manifest.yaml 的 spec 段,四个字段全部必填:
| 字段 | 类型 | 说明 | 默认值 |
|---|---|---|---|
url_region | string | 研究数据所在区域:uk(英国)、nl(荷兰)、us(美国) | uk |
client_id | string | 在对应区域账户设置中生成的 API Client ID | - |
client_secret | string | 对应区域的 API Client Secret | - |
start_date | string | 增量同步的起始时间 | - |
细节说明:
client_id与client_secret均标记airbyte_secret: true,Airbyte 平台会将其作为敏感信息加密存储并在 UI 中脱敏显示;start_date使用format: date-time,并带正则约束^[0-9]{4}-[0-9]{2}-[0-9]{2}T[0-9]{2}:[0-9]{2}:[0-9]{2}Z$,即必须为YYYY-MM-DDTHH:MM:SSZ形式的 UTC 时间戳,与MinMaxDatetime中%Y-%m-%dT%H:%M:%SZ的解析格式严格对齐;- 区域与认证凭据的获取位置,见 Castor 官方用户文档 的 Authentication 一节:
url_region | 服务器 | 账户设置入口 |
|---|---|---|
nl | 荷兰(EU),主机data.castoredc.com | data.castoredc.com/account/settings |
uk | 英国,主机uk.castoredc.com | uk.castoredc.com/account/settings |
us | 美国,主机us.castoredc.com | us.castoredc.com/account/settings |
该文档同时提醒:若使用 Airbyte Cloud 且组织启用了 IP 白名单,需要将 Airbyte Cloud 的出口 IP 加入允许列表,确保连接器能够访问上述三个区域主机。
八、测试与验收:单元测试与 Connector Acceptance Tests
单元测试:直接实例化 manifest 组件
unit_tests/test_manifest.py 是声明式连接器测试的典型范式——不 mock 网络,而是把 manifest 解析为真实 CDK 组件并断言其配置正确性。测试代码通过ModelToComponentFactory().create_component(model_type=HttpRequester, ...)将base_requester反序列化为可执行的HttpRequester对象,然后验证:
- 三个区域的 API 基地址与 OAuth 令牌端点拼接正确;
- 认证类型确为
client_credentials; - 生成的 URL 主机均在
allowedHosts白名单内; url_region的枚举与默认值保持不变。
依赖声明见 unit_tests/pyproject.toml:airbyte-cdk = "7.33.0"(与运行镜像版本一致)、pytest ^8.0、pyyaml ^6.0。
验收测试:CAT 配置
acceptance-test-config.yml 声明了 Connector Acceptance Tests(CAT)的入口,其中spec测试直接指向manifest.yaml(即验证 manifest 生成的 spec 符合协议规范);而connection、discovery、basic_read、incremental、full_refresh均以bypass_reason: "This is a builder contribution, and we do not have secrets at this time"跳过——这是因为该连接器由 Connector Builder 社区贡献,暂无测试凭据。需要说明的是,metadata.yaml的metadata.testedStreams段记录了 16 个流的测试结果标记(hasRecords、primaryKeysAreUnique、responsesAreSuccessful等均为 true),说明这些流的响应结构已通过构建侧校验。
九、本地开发与连接器特定指南
README 的 Development 一节指出:声明式连接器的本地开发与测试流程遵循 Airbyte 的本地连接器开发指南,核心思路是在本地构建镜像后通过spec、check、discover、read等协议命令驱动连接器,并配合单元测试与 CAT 验证行为。对 manifest-only 连接器而言,本地迭代的主要对象就是 manifest.yaml 本身,修改后重新构建airbyte/source-castor-edc:dev镜像即可验证(对应 acceptance-test-config.yml 中的connector_image: airbyte/source-castor-edc:dev)。
README 的 Connector-Specific Guidance 一节还提到:连接器特定的排错与测试指导可以记录在连接器目录下的CONTRIBUTING.md中(当前仓库快照的该目录下未包含此文件,实际以发行产物为准)。这提醒连接器维护者:把"只有本连接器才成立"的踩坑经验与通用模板 README 分开存放,是 Airbyte 连接器仓库的约定做法。
十、总结
source-castor-edc是理解 Airbyte 声明式连接器的最佳范例之一:它在 manifest.yaml 中浓缩了声明式连接器的几乎所有核心能力——基于client_credentials的 OAuth 认证、通过 Jinja 表达式实现的区域动态路由、SubstreamPartitionRouter的父子流分区同步、PageIncrement分页、429 限流重试、DatetimeBasedCursor增量同步,以及AddFields字段变换——并且全部由 单元测试 锁定了关键行为。阅读这类连接器时,正确的姿势是:以 README 为入口,以 manifest 为实现,以 metadata 为注册信息,以测试为行为契约。当你需要为下一个 REST API 编写连接器时,这套"声明式三板斧"完全可以直接复用。
- 数据工程
- 数据集成
- ETL
- 后端
- 大数据
【免费下载链接】airbyte
Open-source data movement for ELT pipelines and AI agents — from APIs, databases & files to warehouses, lakes, and AI applications. Both self-hosted and Cloud.
相关推荐
Airbyte Gutendex 声明式 Source 连接器深度解析:基于 Low-Code CDK 的 manifest 实现与实战配置
Airbyte Gutendex 声明式 Source 连接器深度解析:基于 Low Code CDK 的 manifest 实现与实战配置 本篇文章以 Air
数据工程数据集成ETL后端大数据Airbyte Employment Hero 声明式连接器深度解析:基于 Low-Code CDK 的 manifest 实现与本地开发指南
Airbyte Employment Hero 声明式连接器深度解析:基于 Low Code CDK 的 manifest 实现与本地开发指南 本篇技术指南以开
数据工程数据集成ETL后端大数据Airbyte News API Source 连接器实战指南:manifest-only 声明式连接器的架构、配置与测试
Airbyte News API Source 连接器实战指南:manifest only 声明式连接器的架构、配置与测试 本文以 Airbyte 仓库中的 s
数据工程数据集成ETL后端大数据
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考