数据工程数据集成ETL后端大数据【免费下载链接】airbyteOpen-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点击查看免费下载本文以source-tiktok-marketing连接器专属开发指南CLAUDE.md其内容与 AGENTS.md 相同CLAUDE.md 是指向 AGENTS.md 的符号链接为骨架结合连接器源码manifest.yaml、components.py与单元测试系统讲解该连接器偏离标准声明式declarative连接器模式的六大陷阱动态沙箱/生产端点选择、双广告主 ID 分区路由器、空指标-字符串转换、基于响应体 code 的限流检测、Smart Ads 缺失modify_time过滤以及沙箱账户限流与凭据锁定。读完本文你将理解每个gotcha背后的 API 约束、源码实现位置与影响能够在改动该连接器时避免踩坑。背景这是一个混合式连接器在深入六大行为之前先明确该连接器的技术架构这是理解后续所有内容的前提。source-tiktok-marketing采用声明式清单Declarative Manifest Python 自定义组件Custom Components的混合架构连接器主体由 manifest.yaml约 6200 行version: 1.1.0type: DeclarativeSource声明但其中四个核心自定义组件以 Python 实现于 components.pySingleAdvertiserIdPerPartition单个广告主 ID 分区路由器MultipleAdvertiserIdsPerPartition多个广告主 ID 批量分区路由器TransformEmptyMetrics空指标转换错误处理器中的各 code 谓词同样声明在 manifest 中在 manifest 中这些组件通过class_name: source_declarative_manifest.components.Xxx引用例如 manifest.yaml 第 111 行与第 129 行 引用了两个分区路由器。因此任何针对该连接器的修改都同时涉及 YAML 清单与 Python 代码两个层面。单元测试位于 unit_tests/test_components.py其中覆盖了分区路由器取值test_get_partition_value_from_config、单/多 ID 切片生成test_stream_slices_single/test_stream_slices_multiple以及空指标转换test_transform_empty_metrics可作为验证修改正确性的回归测试入口。一、动态沙箱与生产端点选择1.1 问题描述TikTok Marketing API 存在两套完全不同的 API 基础 URL连接器会根据配置中的认证类型在两者间动态切换环境基础 URL沙箱Sandboxhttps://sandbox-ads.tiktok.com/open_api/v1.3/生产Productionhttps://business-api.tiktok.com/open_api/v1.3/这一选择通过 manifest.yaml 第 19 行 的url_baseJinja 表达式实现url_base: https://{{ sandbox-ads if config.get(credentials, {}).get(auth_type, ) sandbox_access_token else business-api }}.tiktok.com/open_api/v1.3/其逻辑等价于当config.credentials.auth_type sandbox_access_token时走沙箱域名sandbox-ads否则一律走生产域名business-api。也就是说认证方式本身就决定了请求发往哪个环境无需用户单独配置环境开关。1.2 为什么重要沙箱与生产 API 在数据可用性和限流策略上存在显著差异沙箱账户无法通过 API 获取广告主 IDoauth2/advertiser/get/端点在沙箱环境下不工作。这正是配置项允许直接指定advertiser_id的根本原因详见第二节分区路由器如何消费该配置。行为差异使用沙箱账户测试时部分 stream 可能行为不同或返回空数据与生产环境表现不一致。对开发者的启示如果你针对沙箱账户验证修改切不可将沙箱下的表现直接等同于生产行为反之修复了某个沙箱下的问题也必须在生产凭据下做回归验证。二、双广告主 ID 分区路由器2.1 问题描述连接器使用两个自定义分区路由器SubstreamPartitionRouter子类根据advertiser_id是来自配置还是来自 API 父流采用不同处理方式SingleAdvertiserIdPerPartition大多数 stream 使用若配置中存在advertiser_id只产出一个包含该 ID 的分区并完全跳过父流读取否则从advertisers父流advertiser_idsparent stream为每个广告主 ID 各产出一个分区。MultipleAdvertiserIdsPerPartition仅advertisersstream 使用将最多100 个广告主 ID 批量打包进单个 JSON 数组字符串分区例如[id1, id2, ...]因为 TikTok 广告主信息端点支持单次请求携带多个 ID。两个路由器都按优先级顺序检查多个配置路径credentials.advertiser_id→environment.advertiser_id。对应 manifest 中的path_in_config定义manifest.yaml 第 113-115 行path_in_config: - [credentials, advertiser_id] - [environment, advertiser_id]2.2 源码实现两个类都定义在 components.pyclass MultipleAdvertiserIdsPerPartition(SubstreamPartitionRouter): def stream_slices(self) - Iterable[StreamSlice]: partition_value_in_config self.get_partition_value_from_config() if partition_value_in_config: slices [partition_value_in_config] else: slices [_id.partition[self._partition_field] for _id in super().stream_slices()] start, end, step 0, len(slices), 100 for i in range(start, end, step): yield StreamSlice(partition{advertiser_ids: json.dumps(slices[i : min(end, i step)]), parent_slice: {}}, cursor_slice{}) class SingleAdvertiserIdPerPartition(MultipleAdvertiserIdsPerPartition): def stream_slices(self) - Iterable[StreamSlice]: partition_value_in_config self.get_partition_value_from_config() if partition_value_in_config: yield StreamSlice(partition{self._partition_field: partition_value_in_config, parent_slice: {}}, cursor_slice{}) else: yield from super(MultipleAdvertiserIdsPerPartition, self).stream_slices()关键细节get_partition_value_from_config()使用dpath.get(self.config, path, defaultNone)按优先级依次探测两个配置路径返回第一个非空值MultipleAdvertiserIdsPerPartition用json.dumps把 ID 列表序列化为 JSON 数组字符串并按照100 个一批step 100切分分区SingleAdvertiserIdPerPartition继承前者但覆写stream_slices配置存在时直接 yield 单个分区不再调用父流否则回退到父流逐个 ID 产出分区。2.3 为什么重要advertisersstream 以 JSON 数组字符串形式把广告主 ID 放在请求参数中而不是单个值。如果你改动广告主 ID 的分区方式必须牢记advertisersstream 需要批量数组格式其余所有 stream 需要单个 ID。两者一旦混淆后果是向期望数组的端点发送单个 ID 会造成数据缺失向期望单个 ID 的端点发送数组会引发API 错误。这一点在接入oauth2/advertiser/get/不可用的沙箱场景时尤为关键——此时用户必须在配置里显式提供advertiser_id路由器才会跳过父流请求。三、空指标值返回为破折号字符串3.1 问题描述TikTok 报表 API 对没有数据的指标返回字符串-字面破折号而不是null或0。自定义转换TransformEmptyMetrics遍历每条报表记录中的metrics对象把所有-值转换为null。3.2 源码实现components.py 中实现仅十余行dataclass class TransformEmptyMetrics(RecordTransformation): empty_value - def transform(self, record, configNone, stream_stateNone, stream_sliceNone): for metric_key, metric_value in record.get(metrics, {}).items(): if metric_value self.empty_value: record[metrics][metric_key] None return record在 manifest.yaml 中TransformEmptyMetrics被挂载到几乎所有报表 stream的 transformation 链上class_name: source_declarative_manifest.components.TransformEmptyMetrics出现在ads_reports_daily、ads_reports_by_country_daily、ad_groups_reports_daily、advertisers_reports_daily、campaigns_reports_daily、ads_reports_hourly、ads_reports_lifetime及各类 audience / by-country / by-platform / by-province 报表等 30 余处例如 manifest.yaml 第 735 行、第 820 行、第 1552 行。3.3 为什么重要若不经过该转换下游 schema 期望数值类型如spend、clicks、impressions的指标会收到字符串值导致目标端destination类型错误。单元测试test_transform_empty_metricsunit_tests/test_components.py直接验证了这一转换行为。如果你新增报表 stream必须把TransformEmptyMetrics挂入其 transformation 链否则该 stream 会输出非法指标类型。四、通过响应体 code 检测限流4.1 问题描述TikTok 的 API不使用标准 HTTP 429 状态码来表示限流。相反它在HTTP 200的 JSON 响应体里通过code字段报告错误。错误处理器使用以下谓词predicate组合code含义与动作40100限流action: RATE_LIMITED50000/51041/51004/51002瞬时服务端错误action: RETRY自动重试60001服务端维护中action: RETRY40001权限不足action: FAILfailure_type: config_error40002资源不可访问/不存在action: IGNORE40067查询体量超限报表专用action: FAILfailure_type: config_error提示调小 Daily Reports Date Step! 0通用 API 错误action: FAIL全局错误处理器定义于 manifest.yaml 第 22-56 行例如限流检测- predicate: {{ response.get(code) 40100 }} action: RATE_LIMITED error_message: TikTok Marketing API rate limit exceeded. Please verify that only one Airbyte connection with the same credentials is running at a time. ...4.2 为什么重要基于标准 HTTP 状态码的限流检测对 TikTok API 完全无效。若你改动错误处理器必须保留这些响应体 code 检查。限流错误信息特别警告了使用相同凭据的并发连接问题——TikTok 的限流是按访问令牌per-access-token计量的。同时存在一个重复的report_daily_error_handlermanifest.yaml 第 424 行起专门用于日级报表额外把40067查询过大暴露为config_error引导用户调小 Daily Reports Date Step 设置如 7 或 1。全局请求器配置max_retries: 9配合ConstantBackoffStrategy固定 60 秒退避重试节奏较长对慢恢复的限流与维护窗口是必要的。五、Smart Ads 缺失 modify_time 过滤5.1 问题描述adsstream 带有RecordFilter会丢弃modify_time为None的记录。这是专门为处理 TikTokSmart Ad记录而设的该类型广告记录有时会被 API 返回且不带modify_time字段。由于adsstream 以modify_time作为增量游标manifest.yaml 第 99 行cursor_field: modify_time缺失该字段会导致游标比较失败。manifest 中的实现manifest.yaml 第 375-382 行record_selector: $ref: #/definitions/record_selector # This filter is needed because the API will at times return Smart Ad Records without a modify_time value. # These are not easily filtered at the API level which is why they are filtered here. record_filter: type: RecordFilter condition: {{ record.get(modify_time) is not none }}5.2 为什么重要该过滤器会静默丢弃本属有效的广告记录。若用户反馈ads 数据缺失Smart ads 缺少modify_time是首要怀疑对象schema 中modify_time字段的描述也注明 Smart ad 的 ID 仅在 Smart 广告中存在见 manifest.yaml 第 3479 行 附近。这是为维持增量同步可靠性而接受的已知取舍manifest 注释同时给出了上游 PR 上下文链接。六、沙箱账户限流与凭据限制6.1 问题描述TikTok 沙箱账户的限流为10 次请求/秒。如果在 CI 中运行连接器验收测试CATs见 acceptance-test-config.yml的同时又用同一套凭据在本地并行测试就会超过该限流凭据可能被临时限制——导致所有请求失败。已知行为部分来自实践观测TikTok 官方文档没有说明限制大约持续数小时有证据表明在限制期内持续发起请求会延长锁定时长TikTok 官方对该限制行为及其精确时长无文档说明。6.2 为什么重要与大多数排队或重试即可的 API 限流不同超出 TikTok 沙箱限流可能把凭据完全锁死数小时绝不并发不要在 CI 与本地同时针对沙箱账户跑测试。遇锁即停如果沙箱凭据突然出现 100% 请求失败立即停止所有请求并等待不要反复重试。增量流Incremental Stream注意事项TikTok Marketing API 支持基于日期的报表端点过滤。该连接器使用 manifest 引用的 Python 自定义组件来实现增量逻辑。连接器类型Python 自定义组件混合 manifest Python。分析状态stream 通过自定义组件在 Python 侧定义完整的逐 stream 增量分析需要 Python 代码审查。未来的增量 stream 候选所有 stream 均延后至 Python 代码审查本连接器的 stream 定义在 Python 代码中而非纯声明式 manifest YAML。按标准 CONTRIBUTING.md 模板要求应待后续 Agent 审查 Python stream 定义、其cursor_field属性以及它们调用的 API 端点之后再补充完整的逐 stream 增量分析表。从现有 manifest 可以观察到基础流campaigns、ads、ad_groups等通过semi_incremental_syncmanifest.yaml 第 97-108 行实现客户端侧增量——游标为modify_time支持%Y-%m-%d %H:%M:%S与%Y-%m-%dT%H:%M:%SZ两种格式start_date默认2016-09-01而报表类 streamdaily / hourly / lifetime / audience / by-country 等 30 余个见 manifest.yaml 第 647 行起则通过stream_intervalstart/end date做按日切片请求。小结修改本连接器前的检查清单综合以上六点改动source-tiktok-marketing前的自查要点端点url_base的 Jinja 条件sandbox_access_token→ 沙箱不可破坏沙箱下oauth2/advertiser/get/不可用。分区advertisers流用MultipleAdvertiserIdsPerPartition100 个/批、JSON 数组字符串其余流用SingleAdvertiserIdPerPartition单个 ID配置优先级credentials.advertiser_id→environment.advertiser_id。指标所有报表流必须保留TransformEmptyMetrics否则-字符串会破坏数值 schema。错误处理保留基于响应体code40100 限流 / 5xxxx 重试 / 40067 配置错误等的谓词检查勿依赖 HTTP 429。ads 流保留modify_time is not none的RecordFilter理解 Smart 广告会被静默丢弃的取舍。沙箱凭据CI 与本地测试切勿并发使用同一套沙箱凭据遇锁即停。每一条行为都有对应的源码落点manifest.yaml / components.py与测试覆盖unit_tests/test_components.py可作为后续修改与回归验证的直接依据。赞分享数据工程数据集成ETL后端大数据【免费下载链接】airbyteOpen-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点击查看免费下载相关推荐Airbyte 的 source-google-search-console 连接器三大非显而易见行为与限流、兼容性实战解析Airbyte 的 source google search console 连接器三大非显而易见行为与限流、兼容性实战解析 本篇技术指南围绕 Airbyte数据工程数据集成ETL后端大数据Airbyte source-linkedin-ads 连接器深度解析七大非显而易见行为与增量同步设计Airbyte source linkedin ads 连接器深度解析七大非显而易见行为与增量同步设计 导读 本文基于 Airbyte 仓库中 source数据工程数据集成ETL后端大数据Airbyte source-linkedin-ads 连接器深度解析7 大非显而易见行为与增量同步设计Airbyte source linkedin ads 连接器深度解析7 大非显而易见行为与增量同步设计 本篇技术指南以 Airbyte 仓库中 source数据工程数据集成ETL后端大数据创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考