- 数据工程
- 数据集成
- 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.
导读
Apple Search Ads(Apple Ads)是苹果官方的搜索广告投放平台。本文以 Airbyte 仓库中source-apple-search-ads连接器为核心,完整讲解其基于 Airbyte Low-Code CDK 声明式配置实现的数据流设计:实体类 Base Streams(Campaigns / AdGroups / Keywords / Ads)与统计类 Report Streams(四个_daily报告流)的同步模式、增量游标机制、OAuth 认证、重试与分页策略,以及全部配置参数的取值与含义。读完本文,你将掌握该连接器的数据获取边界、同步行为与调优手段,能够据此正确配置源并诊断同步问题。
一、连接器概览:REST API 之上的声明式实现
Apple Search Ads 对外提供的是 REST 风格 API。source-apple-search-ads连接器并非手写 Python/Java 实现,而是采用Airbyte Low-Code CDK的声明式描述方式:所有请求、认证、分页、错误处理与增量逻辑都定义在 manifest.yaml 中,由 CDK 运行时(source-declarative-manifest镜像)解释执行。这一点与仓库中 README.md 描述的 "declarative connector built with the Connector Builder" 一致,其 metadata.yaml 也标注了cdk:low-code与language:manifest-only标签。
API 基址定义在base_requester中:
https://api.searchads.apple.com/api/v5连接器整体将数据流划分为两大类,对应 bootstrap.md 的核心脉络:
| 类别 | 用途 | 同步模式 |
|---|---|---|
| Base streams | API 中的实体属性(有哪些 Campaign、AdGroup、Keyword、Ad) | 仅 Full Refresh |
| Report streams | 实体统计指标(Campaign 花了多少钱、Keyword 有多少次点击等) | Full Refresh + Incremental |
二、Base Streams:实体类数据流
Base streams 返回 Apple Search Ads 账户中实体的属性信息,全部只支持全量刷新(full refresh)。连接器定义了 4 个:
campaigns(推广计划)adgroups(广告组)keywords(关键词)ads(广告)
每个流的实体主键都是id(整数类型),这可以从 configured_catalog.json 中看到:四个流均声明"source_defined_primary_key": [["id"]]与"supported_sync_modes": ["full_refresh"]。
在 manifest.yaml 中,Base streams 的请求路径体现了 Apple Search Ads API 的层级关系:
| 流 | HTTP | 路径 |
|---|---|---|
campaigns | GET | /campaigns |
adgroups | GET | /campaigns/{{ stream_slice.campaign_id }}/adgroups |
keywords | GET | /campaigns/{campaign_id}/adgroups/{adgroup_id}/targetingkeywords |
ads | GET | /campaigns/{campaign_id}/adgroups/{adgroup_id}/ads |
注意后三个流的路径中带有模板变量:AdGroups 隶属于某个 Campaign,Keywords 与 Ads 又嵌套在 AdGroup 之下。这对应 manifest 中的SubstreamPartitionRouter——adgroups以campaigns.id为父键分区,keywords/ads再以adgroups.id为父键二次分区,从而实现"先拉取 Campaign 列表,再逐 Campaign 拉取 AdGroup,再逐 AdGroup 拉取 Keyword / Ad"的层级遍历。
所有 Base stream 请求都会携带组织上下文头X-AP-Context: orgId={{ config.org_id }},即每次调用都必须指明数据所属的组织(Org)。
三、Report Streams:统计类数据流与 DAILY 粒度
Report streams 返回各实体的统计指标(花费、点击、展示等),同时支持全量刷新与增量同步。连接器提供 4 个:
campaigns_report_daily(Campaign 级报告)adgroups_report_daily(Ad Group 级报告)keywords_report_daily(Keyword 级报告)ads_report_daily(Ad 级报告)
当前报告流只设置为DAILY粒度,即请求体中的granularity: DAILY,因此实际数据流名称统一带有_daily后缀。这也是 bootstrap.md 明确指出的现状。
从 configured_catalog.json 可以看到报告流的增量特征:
- 默认游标字段(cursor)为
date,且由源定义("source_defined_cursor": true); - 支持
full_refresh与incremental两种同步模式; - 复合主键为
date+ 实体 ID:campaigns_report_daily为[date, campaignId],adgroups_report_daily为[date, adGroupId],keywords_report_daily为[date, keywordId],ads_report_daily为[date, adId]; - 目标写入模式为
append(追加),增量同步结果按天持续追加到目标端。
3.1 报告接口的调用方式
报告流与实体流不同,使用的是 POST 请求,且分区路由只到 Campaign 层级(按campaign_id分区):
| 流 | HTTP | 路径 |
|---|---|---|
campaigns_report_daily | POST | /reports/campaigns |
adgroups_report_daily | POST | /reports/campaigns/{campaign_id}/adgroups |
keywords_report_daily | POST | /reports/campaigns/{campaign_id}/keywords |
ads_report_daily | POST | /reports/campaigns/{campaign_id}/ads |
每个报告请求的 JSON 请求体(request_body_json)结构如下:
request_body_json: startTime: "{{ stream_slice.start_time }}" endTime: "{{ stream_slice.end_time }}" granularity: DAILY groupBy: "[ 'countryOrRegion' ]" selector: '{ "orderBy": [ { "field": "countryOrRegion", "sortOrder": "ASCENDING" } ] }' timeZone: "{{ config['timezone'] or 'UTC' }}"startTime/endTime来自增量同步产生的日期切片(slice),见下文;granularity固定为DAILY;groupBy按国家/地区(countryOrRegion)分组统计,selector中同时按该字段升序排序;timeZone决定报告统计的时区口径,默认UTC,也可配置为ORTZ(Organization Time Zone,组织时区)。
报告接口的响应被封装在data.reportingDataResponse.row路径下,因此 manifest 中DpathExtractor的field_path为["data", "reportingDataResponse", "row"],分页的 offset/limit 则注入到请求体 JSON 的selector.pagination中(每页 1000 条)。
3.2 报告的字段变换(transformations)
Apple 报告接口返回的原始行数据把指标与元信息混在metadata字段中,因此 manifest 通过AddFields变换把关键信息提升为顶层字段,例如campaigns_report_daily:
transformations: - type: AddFields fields: - type: AddedFieldDefinition path: [campaignId] value: "{{ record.metadata.campaignId }}" - type: AddedFieldDefinition path: [date] value: "{{ stream_slice.start_time }}" - type: AddedFieldDefinition path: [countryorregion] value: "{{ record.metadata.countryOrRegion }}"- 实体 ID(
campaignId/adGroupId/keywordId/adId)从record.metadata中提出; date直接取当前日期切片起点stream_slice.start_time;countryorregion同样从metadata.countryOrRegion提取。
将countryOrRegion单独提取为顶层字段是有实际意义的:报告按国家/地区分组后,若仍以整个metadata字段作为主键进行去重,metadata中其他键的变化会导致同一date+ 实体 ID 的记录出现重复;单独提取后即可用countryOrRegion代替整个metadata参与主键去重。这也是 manifest.yaml 顶部 description 中解释的修复动机。
四、增量同步机制:start_date、end_date 与日粒度切片
报告流支持增量同步,其核心是 manifest 中的DatetimeBasedCursor(时间游标)。bootstrap.md 明确说明:
连接器使用
start_date配置项作为首次报告同步的起点;若未显式设置end_date,则以当前日期作为结束日期。
对应 manifest 的实现为:
incremental_sync: type: DatetimeBasedCursor cursor_field: date lookback_window: P{{ config.lookback_window }}D cursor_datetime_formats: ["%Y-%m-%d"] datetime_format: "%Y-%m-%d" start_datetime: type: MinMaxDatetime datetime: "{{ config.start_date }}" datetime_format: "%Y-%m-%d" end_datetime: type: MinMaxDatetime datetime: "{{ config.end_date or today_utc() }}" datetime_format: "%Y-%m-%d" step: P1D cursor_granularity: P1D关键点:
- 日期格式:全部为
%Y-%m-%d(如2022-11-11); - 起点:
start_date(配置必填); - 终点:
config.end_date or today_utc()——配置了end_date用配置值(含当天,文档描述为 "Data is retrieved until that date (included)"),否则用运行当天 UTC 日期; - 切片步长:
step: P1D表示按"每天"切分请求区间,即每个 slice 覆盖一天,对应 DAILY 粒度; - 回看窗口:
lookback_window: P{{ config.lookback_window }}D让每次增量都额外向前回看若干天,以吸收 Apple 归因数据的延迟更新(见配置参数表)。
4.1 从分区状态到全局游标:1.0.0 破坏性变更
从源码结构看,adgroups_report_daily、keywords_report_daily、ads_report_daily三个流都声明了global_substream_cursor: true。这与 metadata.yaml 中记录的 1.0.0 破坏性变更直接对应:
该版本将
adgroups_report_daily与keywords_report_daily的状态从"按分区(per-partition)状态"改为"使用全局状态游标(global state cursor)",从而缩短这两个流的读取时间。
campaigns_report_daily未在变更影响范围内,说明其从更早版本起即使用全局游标。对于从旧版本升级的用户,需要清空(clear)受影响流的历史数据后再同步,迁移截止时间为 2025-11-04。这一升级细节是排查"升级后状态不兼容"问题的重要线索。
五、配置参数详解
连接器的spec定义在 manifest.yaml 的spec.connection_specification中。必填项为:org_id、client_id、start_date、client_secret、timezone、token_refresh_endpoint。可参考 sample_config.json 中的示例结构:
{ "org_id": "REPLACEME", "client_id": "REPLACEME", "client_secret": "REPLACEME", "token_refresh_endpoint": "https://apple.oauth.com/token", "start_date": "2022-11-11", "backoff_factor": 10, "lookback_window": 3 }| 参数 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
org_id | integer | 是 | — | 拥有 Campaign 的组织标识符,与 Apple Search Ads 控制台中的账户(Org)一致,随请求头X-AP-Context发送 |
client_id | string | 是 | — | 获取令牌所用的用户标识(OAuth client id),secret 类型 |
client_secret | string | 是 | — | 认证用户设置请求的客户端密钥,secret 类型 |
start_date | string | 是 | — | 首次同步数据的起始日期,格式YYYY-MM-DD(正则^[0-9]{4}-[0-9]{2}-[0-9]{2}$) |
end_date | string | 否 | 当前日期 | 数据检索截至日期(含当天),格式同start_date;未设置则以运行当天 UTC 为终点 |
timezone | string | 是 | UTC | 报告统计时区,仅UTC或ORTZ(组织时区) |
token_refresh_endpoint | string | 是 | https://appleid.apple.com/auth/oauth2/token?grant_type=client_credentials&scope=searchadsorg | OAuth 令牌刷新端点,需要代理 Apple 令牌请求时可覆盖 |
backoff_factor | integer | 否 | 5 | 指数退避的延迟增长系数,有效值 1–20(正则^(20|1[0-9]|[1-9])$) |
lookback_window | integer | 否 | 30 | 增量同步回看天数(Apple 采用 30 天归因窗口;调小可缩短同步耗时,代价是可能漏掉延迟归因数据,示例值为 7) |
num_workers | integer | 否 | 2 | 同步并发工作线程数,范围 1–20,配合concurrency_level(默认取num_workers,上限 20)使用 |
其中start_date/end_date的正则约束、num_workers的 1–20 范围、backoff_factor与lookback_window的取值范围都能在 manifest 的 spec 中直接找到依据。sample_config.json中的token_refresh_endpoint是测试环境端点,生产环境应按需使用默认的 Apple 端点。
六、认证与错误处理
6.1 OAuth client_credentials
所有请求都通过OAuthAuthenticator完成认证(manifest 的base_requester):
authenticator: type: OAuthAuthenticator client_id: "{{ config.client_id }}" grant_type: client_credentials client_secret: "{{ config.client_secret }}" token_refresh_endpoint: "{{ config.get('token_refresh_endpoint', 'https://appleid.apple.com/auth/oauth2/token?grant_type=client_credentials&scope=searchadsorg') }}" refresh_request_body: {}即使用client_credentials授权模式,用client_id+client_secret换取访问令牌;端点可在token_refresh_endpoint中显式覆盖。
6.2 三级重试策略
manifest 为每个流都配置了CompositeErrorHandler/DefaultErrorHandler配合HttpResponseFilter:
- 401 → 刷新令牌后重试:
action: REFRESH_TOKEN_THEN_RETRY,failure_type: transient_error,错误信息为 "Access token is expired."——当 Apple 在 CDK 记录的令牌过期时间之前就返回 401 时,连接器会先主动刷新令牌再重试; - 500 / 429 → 直接重试:
action: RETRY,其中 429 是限流、500 是服务端错误; - 指数退避:
ExponentialBackoffStrategy按config.backoff_factor控制延迟增长,报告流(max_retries: 10)的重试上限高于实体流。
此外,keywords_report_daily有一个特殊过滤规则:当错误消息包含CAMPAIGN DOES NOT CONTAIN KEYWORD时执行IGNORE——即"该 Campaign 不含关键词"并非异常,直接跳过该切片,避免因个别空数据分区导致整个同步失败。
这些行为已被 unit_tests/test_manifest.py 固化为断言:如test_streams_reactively_refresh_oauth_token_on_401校验每个流都具备 401 刷新重试、test_keywords_report_daily_retains_keyword_predicate校验关键词流的 IGNORE 谓词、test_ads_report_daily_no_keyword_error_predicate则防止该谓词被错误地复制到 ads 报告流。
七、分页与并发
- 实体流分页:采用
OffsetIncrement策略,通过查询参数offset与limit翻页,每页page_size: 1000; - 报告流分页:同样的 offset/limit 策略,但注入位置是请求体 JSON 的
selector.pagination.offset/selector.pagination.limit; - 并发控制:manifest 顶层声明
ConcurrencyLevel,默认并发数为config.get('num_workers', 2),上限 20。num_workers默认 2、范围 1–20。由于 Keywords / Ads 属于两层嵌套的子流(先遍历 Campaign 再遍历 AdGroup),并发分区处理能显著缩短深层嵌套流的同步时间——test_concurrency_level_configured明确注释了这是为了"防止深层嵌套子流在心跳超时内无法完成"。
八、测试与质量保障
连接器的测试配置展示了针对该数据源的实际验证手段:
- unit_tests/test_manifest.py:以数据驱动方式校验 manifest 的日期字段变换(
date必须取自stream_slice.start_time而非stream_slice.start_date)、401 刷新、关键词谓词、并发配置与num_workersspec 字段; - acceptance-test-config.yml:声明
spec与connection两类验收测试,其中连接测试使用无效配置(invalid_config.json 包含wrongid、非法日期9999-99-99、非法backoff_factor等)验证校验器能正确拒绝;full_refresh 验收测试对granularity与metadata字段做了"天然不可幂等"的豁免说明; - integration_tests/sample_state.json 与 abnormal_state.json:分别给出正常游标(
date: 2022-11-09)与未来游标(date: 2999-11-06),用于验证增量同步在正常与超前状态下的行为。
九、小结:同步行为速查
| 数据流 | 同步模式 | 游标 | 主键 | 请求方式 |
|---|---|---|---|---|
campaigns | full_refresh | — | id | GET/campaigns |
adgroups | full_refresh | — | id | GET/campaigns/{id}/adgroups |
keywords | full_refresh | — | id | GET/campaigns/{cid}/adgroups/{aid}/targetingkeywords |
ads | full_refresh | — | id | GET/campaigns/{cid}/adgroups/{aid}/ads |
campaigns_report_daily | full_refresh + incremental | date | date + campaignId | POST/reports/campaigns |
adgroups_report_daily | full_refresh + incremental | date(全局游标) | date + adGroupId | POST/reports/campaigns/{id}/adgroups |
keywords_report_daily | full_refresh + incremental | date(全局游标) | date + keywordId | POST/reports/campaigns/{id}/keywords |
ads_report_daily | full_refresh + incremental | date(全局游标) | date + adId | POST/reports/campaigns/{id}/ads |
实际使用时,实体类数据建议配合报告流一起同步:实体流提供 Campaign / AdGroup / Keyword / Ad 的元数据与 ID 映射,报告流按 DAILY 粒度持续增量拉取统计指标;start_date决定初始回填范围,end_date不设则默认到当天,lookback_window用于兜底 Apple 延迟归因的数据修正。深入阅读 bootstrap.md 与 manifest.yaml 可获得与本文完全一致的权威细节。
- 数据工程
- 数据集成
- 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 Bing Ads Source 连接器深度解析:从 OAuth 认证、账户分层流到报告与 Bulk 异步下载的完整实现指南
Airbyte Bing Ads Source 连接器深度解析:从 OAuth 认证、账户分层流到报告与 Bulk 异步下载的完整实现指南 本篇技术指南以 Ai
数据工程数据集成ETL后端大数据Airbyte source-amazon-ads 连接器深度剖析:异步报告生成、HTTP 425 冲突与增量同步的独特行为
Airbyte source amazon ads 连接器深度剖析:异步报告生成、HTTP 425 冲突与增量同步的独特行为 本篇技术指南围绕 Airbyte
数据工程数据集成ETL后端大数据Airbyte source-bing-ads 连接器独特行为深度解析:谓词去重过滤器与 Bulk 报告 gzip 解码回退机制
Airbyte source bing ads 连接器独特行为深度解析:谓词去重过滤器与 Bulk 报告 gzip 解码回退机制 本指南以 source bin
数据工程数据集成ETL后端大数据
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考