Airbyte Apple Search Ads Source 连接器完全指南:Base Streams 与 DAILY 粒度报告流
2026/9/23 12:31:18 网站建设 项目流程
  • 数据工程
  • 数据集成
  • 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.

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

导读

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-codelanguage:manifest-only标签。

API 基址定义在base_requester中:

https://api.searchads.apple.com/api/v5

连接器整体将数据流划分为两大类,对应 bootstrap.md 的核心脉络:

类别用途同步模式
Base streamsAPI 中的实体属性(有哪些 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路径
campaignsGET/campaigns
adgroupsGET/campaigns/{{ stream_slice.campaign_id }}/adgroups
keywordsGET/campaigns/{campaign_id}/adgroups/{adgroup_id}/targetingkeywords
adsGET/campaigns/{campaign_id}/adgroups/{adgroup_id}/ads

注意后三个流的路径中带有模板变量:AdGroups 隶属于某个 Campaign,Keywords 与 Ads 又嵌套在 AdGroup 之下。这对应 manifest 中的SubstreamPartitionRouter——adgroupscampaigns.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_refreshincremental两种同步模式;
  • 复合主键为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_dailyPOST/reports/campaigns
adgroups_report_dailyPOST/reports/campaigns/{campaign_id}/adgroups
keywords_report_dailyPOST/reports/campaigns/{campaign_id}/keywords
ads_report_dailyPOST/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 中DpathExtractorfield_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_dailykeywords_report_dailyads_report_daily三个流都声明了global_substream_cursor: true。这与 metadata.yaml 中记录的 1.0.0 破坏性变更直接对应:

该版本将adgroups_report_dailykeywords_report_daily的状态从"按分区(per-partition)状态"改为"使用全局状态游标(global state cursor)",从而缩短这两个流的读取时间。

campaigns_report_daily未在变更影响范围内,说明其从更早版本起即使用全局游标。对于从旧版本升级的用户,需要清空(clear)受影响流的历史数据后再同步,迁移截止时间为 2025-11-04。这一升级细节是排查"升级后状态不兼容"问题的重要线索。

五、配置参数详解

连接器的spec定义在 manifest.yaml 的spec.connection_specification中。必填项为:org_idclient_idstart_dateclient_secrettimezonetoken_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_idinteger拥有 Campaign 的组织标识符,与 Apple Search Ads 控制台中的账户(Org)一致,随请求头X-AP-Context发送
client_idstring获取令牌所用的用户标识(OAuth client id),secret 类型
client_secretstring认证用户设置请求的客户端密钥,secret 类型
start_datestring首次同步数据的起始日期,格式YYYY-MM-DD(正则^[0-9]{4}-[0-9]{2}-[0-9]{2}$
end_datestring当前日期数据检索截至日期(含当天),格式同start_date;未设置则以运行当天 UTC 为终点
timezonestringUTC报告统计时区,仅UTCORTZ(组织时区)
token_refresh_endpointstringhttps://appleid.apple.com/auth/oauth2/token?grant_type=client_credentials&scope=searchadsorgOAuth 令牌刷新端点,需要代理 Apple 令牌请求时可覆盖
backoff_factorinteger5指数退避的延迟增长系数,有效值 1–20(正则^(20|1[0-9]|[1-9])$
lookback_windowinteger30增量同步回看天数(Apple 采用 30 天归因窗口;调小可缩短同步耗时,代价是可能漏掉延迟归因数据,示例值为 7)
num_workersinteger2同步并发工作线程数,范围 1–20,配合concurrency_level(默认取num_workers,上限 20)使用

其中start_date/end_date的正则约束、num_workers的 1–20 范围、backoff_factorlookback_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

  1. 401 → 刷新令牌后重试action: REFRESH_TOKEN_THEN_RETRYfailure_type: transient_error,错误信息为 "Access token is expired."——当 Apple 在 CDK 记录的令牌过期时间之前就返回 401 时,连接器会先主动刷新令牌再重试;
  2. 500 / 429 → 直接重试action: RETRY,其中 429 是限流、500 是服务端错误;
  3. 指数退避ExponentialBackoffStrategyconfig.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策略,通过查询参数offsetlimit翻页,每页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:声明specconnection两类验收测试,其中连接测试使用无效配置(invalid_config.json 包含wrongid、非法日期9999-99-99、非法backoff_factor等)验证校验器能正确拒绝;full_refresh 验收测试对granularitymetadata字段做了"天然不可幂等"的豁免说明;
  • integration_tests/sample_state.json 与 abnormal_state.json:分别给出正常游标(date: 2022-11-09)与未来游标(date: 2999-11-06),用于验证增量同步在正常与超前状态下的行为。

九、小结:同步行为速查

数据流同步模式游标主键请求方式
campaignsfull_refreshidGET/campaigns
adgroupsfull_refreshidGET/campaigns/{id}/adgroups
keywordsfull_refreshidGET/campaigns/{cid}/adgroups/{aid}/targetingkeywords
adsfull_refreshidGET/campaigns/{cid}/adgroups/{aid}/ads
campaigns_report_dailyfull_refresh + incrementaldatedate + campaignIdPOST/reports/campaigns
adgroups_report_dailyfull_refresh + incrementaldate(全局游标)date + adGroupIdPOST/reports/campaigns/{id}/adgroups
keywords_report_dailyfull_refresh + incrementaldate(全局游标)date + keywordIdPOST/reports/campaigns/{id}/keywords
ads_report_dailyfull_refresh + incrementaldate(全局游标)date + adIdPOST/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.

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

相关推荐

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

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

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

立即咨询