- 数据工程
- 数据集成
- 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.
ChartMogul 是面向订阅型 SaaS 业务的收入数据分析平台,提供客户、MRR/ARR、活动事件等核心指标 API。本篇文章以 Airbyte 仓库中 source-chartmogul 连接器 为核心,深入讲解这个基于 Connector Builder / Low-Code CDK 构建的声明式连接器:如何理解其 manifest 配置、六大数据流的认证与分页机制、连接参数如何定义,以及如何通过 Connector Acceptance Tests 进行本地开发与验证。读完本文,你将掌握阅读和扩展任意 Airbyte 声明式连接器所需的核心技能,并能够直接上手使用 ChartMogul 连接器同步订阅业务数据。
一、连接器定位:一个纯声明式(Manifest-Only)的 Low-Code 连接器
在 source-chartmogul/README.md 开头明确说明:这是一个使用Connector Builder构建的声明式连接器(declarative connector),其底层格式遵循Low-Code CDK(即 Config-Based CDK)的 YAML 规范。与传统的 Python/Java 手写连接器不同,这类连接器的全部逻辑——请求构造、认证、记录提取、分页、Schema——都通过一份描述性 YAML 清单(manifest)声明出来,无需编写任何运行时代码。
从仓库文件结构看,该连接器目录下没有source.py、main.py之类的 Python 实现文件,取而代之的是:
- manifest.yaml(1224 行):连接器的唯一"源代码",定义了全部数据流、认证方式、分页策略与 Schema;
- metadata.yaml:连接器元数据(版本、发布状态、仓库信息、破坏性变更说明等);
- acceptance-test-config.yml:Connector Acceptance Tests(CAT)测试套件配置;
- integration_tests/:示例配置、预期记录、测试目录。
这一点在 metadata.yaml 的tags中得到印证:cdk:low-code与language:manifest-only。该连接器镜像为airbyte/source-chartmogul,当前版本1.1.49,releaseStage 为beta,supportLevel 为community,license 为 ELv2。
二、连接配置:API Key 与 Start Date
声明式连接器的输入参数定义在 manifest 的spec.connection_specification中。ChartMogul 连接器只需要两个字段(见 manifest.yaml):
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
api_key | string | 是 | ChartMogul API Key,airbyte_secret: true标记为机密字段(界面掩码显示),order: 0决定表单展示顺序 |
start_date | string | 是 | UTC 格式的起始时间,如2017-01-25T00:00:00Z,格式受pattern: ^[0-9]{4}-[0-9]{2}-[0-9]{2}T[0-9]{2}:[0-9]{2}:[0-9]{2}Z$约束,format: date-time,order: 1 |
一个合法的连接配置 JSON 形如 integration_tests/sample_config.json:
{ "api_key": "<api-key>", "start_date": "2022-01-05T12:09:00Z" }而 integration_tests/invalid_config.json 则用<invalid_key>故意构造失败用例,用于 CAT 中的connection失败测试。start_date的实际作用是数据回填边界:activities流和customer_*_count系列流都会把它作为请求参数传给 ChartMogul API,凡早于该日期的数据不会被同步。
三、数据流全景:6 个 Stream 的声明式定义
整个连接器定义了 6 个数据流(streams),其中 4 个来自同一个 ChartMogul 端点/v1/metrics/customer-count,只是interval粒度不同:
| Stream | API 端点 | 主键 | 分页方式 | 请求参数 |
|---|---|---|---|---|
customers | GET /v1/customers | id | PageIncrement(page/per_page) | 无 |
activities | GET /v1/activities | uuid | CursorPagination(start-after/per_page) | start-date |
customer_daily_count | GET /v1/metrics/customer-count | date | 无 | start-date、end-date、interval=day |
customer_weekly_count | GET /v1/metrics/customer-count | date | 无 | interval=week |
customer_monthly_count | GET /v1/metrics/customer-count | date | 无 | interval=month |
customer_quarterly_count | GET /v1/metrics/customer-count | date | 无 | interval=quarter |
3.1 统一认证:Basic HTTP Auth
所有流共享同一套认证方式(manifest.yaml):
base_requester: type: HttpRequester url_base: https://api.chartmogul.com authenticator: type: BasicHttpAuthenticator username: "{{ config['api_key'] }}" password: "{{ config['api_key'] }}"ChartMogul 的 API 采用 HTTP Basic 认证,且要求用户名和密码都填 API Key 本身。manifest 通过{{ config['api_key'] }}这种 Jinja 模板语法引用运行时配置,CDK 会在请求发出前自动为每个请求附加Authorization: Basic ...头。allowedHosts在 metadata.yaml 中声明为api.chartmogul.com,确保连接器只与官方 API 通信。
3.2 customers:基于页码递增的分页
customers流的主键是id,请求路径/v1/customers,分页采用PageIncrement策略(manifest.yaml):
paginator: type: DefaultPaginator page_token_option: type: RequestOption inject_into: request_parameter field_name: page page_size_option: type: RequestOption inject_into: request_parameter field_name: per_page pagination_strategy: type: PageIncrement start_from_page: 1 page_size: 200即从第 1 页开始,每页 200 条,把页码写入page查询参数,把页大小写入per_page查询参数,逐页递增直到服务器返回空页为止。这适用于数据量相对可控、API 支持页码偏移的端点。
3.3 activities:基于游标的分页
activities流的主键是uuid,请求路径/v1/activities,并通过request_parameters把配置中的start_date原样传入start-date参数。它的分页更精细,采用CursorPagination(manifest.yaml):
pagination_strategy: type: CursorPagination page_size: 200 cursor_value: "{{ response['entries'][-1]['uuid'] }}" stop_condition: "{{ not response.has_more }}"每次请求后,CDK 从响应体entries数组的最后一个元素的uuid字段取出游标值,写入下一请求的start-after参数;只有当响应中的has_more为假时停止翻页。游标分页比页码分页更稳健,能避免在同步过程中数据增删导致的重复或遗漏,适合持续追加的事件类数据。
3.4 customer_*_count 系列:指标聚合流
这四个流(daily/weekly/monthly/quarterly)都请求GET /v1/metrics/customer-count,区别仅在interval参数(day/week/month/quarter),且都通过模板表达式动态计算时间范围(manifest.yaml):
request_parameters: start-date: "{{ format_datetime(config['start_date'], '%Y-%m-%d') }}" end-date: "{{ now_utc().strftime('%Y-%m-%d') }}" interval: daystart-date取自配置但先经format_datetime规范化为YYYY-MM-DD格式,end-date使用 CDK 内置的now_utc()取当前 UTC 时间并格式化为日期,实现"从配置起始日到今天"的全量区间拉取。四个流的主键都是date,Schema 也极简——只有date(string)和customers(integer)两个字段。
需要特别留意的是 metadata.yaml 中记录的破坏性变更:1.0.0 版本把原来的customer_count单一流拆分成了 daily/weekly/monthly/quarterly 四个流。因此旧版本用户升级后需要执行一次 Reset 才能让新的流生效并继续同步。
四、记录提取与 Schema 定义
4.1 记录提取:DpathExtractor
三个不同形态的端点(customers、activities、metrics)返回的数据结构都是"外层包一个entries数组",因此所有流统一使用DpathExtractor并指向entries(manifest.yaml):
record_selector: type: RecordSelector extractor: type: DpathExtractor field_path: - entriesCDK 会从 JSON 响应中按 JSONPath 取出entries数组,逐条作为 Airbyte record 输出。
4.2 连接器级健康检查
连接器的check(连通性验证)声明为CheckStream,指向customers流(manifest.yaml)。即平台在测试连接时,会实际请求一次/v1/customers,只要能正常返回记录(哪怕为空)即判定凭据有效——这也是invalid_config.json中错误 API Key 会导致 CAT 连接测试失败的原因。
4.3 字段 Schema
manifest 中customers流的 Schema(manifest.yaml)完整映射了 ChartMogul 客户对象的字段,包括:
- 身份字段:
id(integer,主键)、uuid(string)、external_id/external_ids(string/array)、data_source_uuid/data_source_uuids、email、name、company; - 状态与时间:
status、state、customer_since、lead_created_at、free_trial_started_at; - 地域信息:
country、state、city、zip及嵌套对象address(含address_zip、city、country、state); - 财务指标:
mrr、arr(integer,注意 manifest 中arr键名与预期记录一致)、currency、currency-sign; - 扩展属性:
attributes(内含clearbit、custom、stripe、tags等子对象)、billing-system-type、billing-system-url、chartmogul-url。
activities流的 Schema(manifest.yaml)则覆盖了活动事件的典型字段:uuid(主键)、date、type、description、currency、activity-mrr、activity-arr、activity-mrr-movement、subscription-external-id、plan-external-id、customer-name/customer-uuid/customer-external-id、billing-connector-uuid。这些字段与 expected_records.jsonl 中的真实样例记录一一对应(如new_biz类型的活动记录包含activity-mrr-movement: 4100、activity-arr: 49200)。
所有 Schema 都设置了additionalProperties: true,并且metadata.autoImportSchema对 6 个流全部显式关闭(manifest.yaml),意味着连接器不依赖 API 自动导入 Schema,而是完全以 manifest 内联定义为准,保证字段结构的稳定性。
五、本地开发与测试实践
原文档指出,本地开发与测试请参照"Developing Connectors Locally"流程;对于声明式连接器而言,改动的核心就是 manifest 文件本身。仓库为该连接器准备了一整套可复用的测试资产:
5.1 Connector Acceptance Tests 配置
acceptance-test-config.yml 声明了 5 类测试套件:
spec:校验 manifest 生成的连接器规范,spec_path: "manifest.yaml";connection:分别用有效配置(secrets/config.json,期望succeed)和无效配置(integration_tests/invalid_config.json,期望failed)验证连通性检查;discovery:执行 Schema 发现;basic_read:按 configured_catalog.json 读取 6 个流,并与expected_records.jsonl比对(exact_order: no,不要求顺序完全一致,fail_on_extra_columns: false容忍额外列);full_refresh:验证全量刷新同步模式下每个流都返回相同记录。
configured_catalog.json 显示所有 6 个流目前仅支持full_refresh同步模式(supported_sync_modes: ["full_refresh"]),主键由 source 定义(source_defined_primary_key),目标侧使用overwrite。integration_tests/acceptance.py则是一个极简的 pytest 插件入口,仅为 CAT 挂载connector_acceptance_test.plugin并预留一个空的connector_setupfixture。
5.2 运行方式
在本地开发时,典型操作是把仓库构建出airbyte/source-chartmogul:dev镜像(CAT 配置中connector_image即指向该镜像),将真实凭据放到secrets/config.json(与sample_config.json结构一致),然后运行:
# 运行连接器镜像的检查命令(示例) docker run --rm -v $(pwd)/secrets:/secrets airbyte/source-chartmogul:dev check --config /secrets/config.json之后通过 Connector Acceptance Tests 依次跑spec、connection、discovery、basic_read、full_refresh五个套件,即可验证 manifest 改动的正确性。需要说明的是:连接器的详细用户文档与逐步设置指南发布在官方文档站点;仓库内若需要补充连接器特有的排障与测试指引,则按原文档约定应追加到该连接器目录下的CONTRIBUTING.md(当前仓库此目录中尚未包含该文件,可按需新增)。
六、总结:声明式连接器带来的工程收益
通过 ChartMogul 这个实例可以清晰看到 Airbyte Low-Code CDK 的工程模式:
- 零代码交付:1224 行的 manifest.yaml 同时承载了连接器规范、认证、6 个流、分页与 Schema,任何具备 YAML 基础的人都可以读懂并修改,无需编译;
- 声明式表达复杂逻辑:
BasicHttpAuthenticator表达 Basic 认证、PageIncrement与CursorPagination表达两种分页语义、{{ ... }}模板表达动态参数,全部是可复用、可组合的 CDK 构件; - 测试资产完备:配合 acceptance-test-config.yml 与 integration_tests/ 目录,任何 manifest 改动都能被 CAT 自动化回归验证。
对于需要在 Airbyte 中同步 ChartMogul 订阅收入数据的团队,本连接器开箱即用:配置好 API Key 与起始日期,即可把客户档案、订阅活动事件和按日/周/月/季粒度的客户数指标持续搬运到你的数据仓库或数据湖中。
- 数据工程
- 数据集成
- 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.
相关推荐
Symfony Dependency Injection组件指南
Symfony Dependency Injection组件指南 还在为PHP应用中的对象依赖管理而头疼吗?每次修改构造函数参数都要到处修改依赖代码?Symfo
数据工程数据集成ETL后端大数据Airbyte Criteo Marketing 声明式源详解:Low-Code CDK 清单、OAuth 认证与增量同步实现
Airbyte Criteo Marketing 声明式源详解:Low Code CDK 清单、OAuth 认证与增量同步实现 Airbyte 中的 sourc
数据工程数据集成ETL后端大数据Airbyte CDK深度解析:构建自定义连接器
Airbyte CDK深度解析:构建自定义连接器 本文深入解析Airbyte CDK架构设计与开发实践,全面对比Python CDK与Java CDK的技术特性
数据工程数据集成ETL后端大数据
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考