- 数据工程
- 数据编排
- ETL
- 任务调度
- 批处理
- 流处理
- 数据集成
- 后端
【免费下载链接】mage-ai
🧙 Build, run, and manage data pipelines for integrating and transforming data.
LinkedIn Ads(领英广告)是营销数据集成场景中的高频数据源。本文以 mage-ai 开源仓库中的 linkedin_ads 数据源模块 为主线,完整讲解该数据源的配置项、两种 OAuth 认证方式、8 个可用数据流(stream)以及底层增量同步、分页与错误重试机制。读完本文,你将能够在 Mage 项目中独立配置并跑通 LinkedIn Ads → 数据仓库/数湖的营销数据同步管道,并理解配置项背后的源码实现。
一、数据源概述
linkedin_ads是 Mage 数据集成框架(mage_integrations)中的一个 Singer 风格数据源(tap),对应的主类为LinkedinAds。它封装了 LinkedIn Marketing Developer Platform 的 Ads API(https://api.linkedin.com/v2),对外暴露三个标准能力:
discover(...):调用 discover.py 读取本地 schema 并生成可选的 catalog 流;sync(...):调用 sync.py 执行真实的数据拉取;test_connection():调用client.check_accounts(config)校验配置中提供的广告账号是否有效。
从源码结构看,LinkedinAds直接继承mage_integrations.sources.base.Source,因此它可以被 Mage 的数据集成管道(data integration pipeline)直接使用,并复用框架提供的 catalog、schema 与状态管理能力。
二、配置参数详解
按照 README.md 的说明,配置该数据源时必须填写以下凭据字段。完整的模板文件位于 templates/config.json,可直接作为配置起点:
{ "client_id": "", "client_secret": "", "refresh_token": "", "accounts": "", "request_timeout": 300, "start_date": "2023-01-01T00:00:00Z", "user_agent": "" }2.1 基础字段
| Key | 说明 | 示例值 | 是否必填 |
|---|---|---|---|
accounts | 需要同步的 LinkedIn 广告账号 ID 列表,逗号分隔;仅当你想同步accounts或account_users流时需要 | "id1, id2, id3" | 同步accounts/account_users时必填 |
access_token | 长期访问令牌 | def789... | ✅ |
request_timeout | 单次 API 请求超时时间(秒) | 300 | 选填 |
start_date | 增量同步的起始时间,ISO 8601 格式 | 2023-01-01T00:00:00Z | ✅ |
user_agent | 请求头中的 User-Agent,建议填联系邮箱 | your_email@your_domain.com | ✅ |
其中request_timeout与源码中的常量REQUEST_TIMEOUT = 300(client.py)保持一致。在 LinkedinClient 构造函数 中可以看到其解析逻辑:
- 当
request_timeout传入的值为非 0 的数值时,转换为float作为真实超时时间; - 当值为
0、"0"或空字符串时,回退到默认的 300 秒。
注意:
accounts字段在同步时会被按逗号拆分并去除空格(config['accounts'].replace(" ", "").split(",")),见 sync.py,因此示例值中的空格是可容忍的。
2.2 替代认证字段(OAuth 三件套)
除了直接提供access_token,README 还给出了另一种更推荐的长效认证方式——提供 OAuth 应用的客户端凭据与刷新令牌:
| Key | 说明 | 示例值 | 是否必填 |
|---|---|---|---|
client_id | LinkedIn 应用的客户端 ID | abc123... | ✅ |
client_secret | LinkedIn 应用的客户端密钥 | xyz456... | ✅ |
refresh_token | 刷新令牌,用于自动续期 access_token | def789... | ✅ |
这两种认证方式在 LinkedinClient 中统一处理:当未提供refresh_token时,视为“旧连接”,直接信任用户传入的access_token;当提供了refresh_token时,客户端会在进入上下文(__enter__)时自动调用fetch_and_set_access_token()判断令牌是否过期并刷新。
2.3 如何获取 access_token(官方流程)
按 README 的步骤指引,获取access_token的完整流程为:
- 登录 LinkedIn 开发者平台,创建一个 LinkedIn 应用;
- 在应用中启用Marketing Developer Platform产品(该产品需要单独申请);
- 填写接入申请表并提交,等待数个工作日的审核批准;
- 审批通过后,使用开发者平台的 OAuth 工具按指引生成 access token。
提示:由于 Marketing Developer Platform 属于受限产品,审核通常需要数天。在等待期内可以先在 Mage 中配置
client_id、client_secret、refresh_token三件套并接入同步逻辑,待令牌可用后再跑通全流程。
三、令牌生命周期管理:从源码看自动续期
linkedin_ads数据源的核心健壮性设计集中在 client.py 的令牌管理逻辑中,这也是它区别于“一次性 access_token”方案的关键。
3.1 令牌刷新流程
客户端维护了三个 OAuth 相关端点常量:
BASE_URL = 'https://api.linkedin.com/v2':Ads API 根地址;LINKEDIN_TOKEN_URI = 'https://www.linkedin.com/oauth/v2/accessToken':令牌刷新端点;INTROSPECTION_URI = 'https://www.linkedin.com/oauth/v2/introspectToken':令牌校验端点。
fetch_and_set_access_token()(client.py)的执行逻辑为:
- 若未配置
refresh_token,直接返回(视为已提供有效 access_token 的旧连接); - 若已配置 access_token,则调用
get_token_expires()调用 introspect 接口获取令牌过期时间; - 若
expires_at晚于当前时间,则日志记录“令牌仍有效”并复用现有令牌; - 否则调用
refresh_access_token(),向accessToken端点提交grant_type=refresh_token换取新的 access_token,并按返回的expires_in(秒)推算新的过期时间。
令牌刷新与校验方法均带有@backoff.on_exception指数退避重试(max_tries=5, factor=2),针对Server5xxError与LinkedInUnauthorizedError自动重试。
3.2 请求统一入口与重试策略
所有 API 请求统一走request()方法(client.py),它会在每次请求前再次调用fetch_and_set_access_token(),确保令牌始终有效;随后自动注入Authorization: Bearer <token>与Accept: application/json请求头,POST 请求额外注入Content-Type。
针对不同失败类型,重试策略分为两套:
- 对 5xx 服务端错误、连接错误与 429 限流,使用
max_time=600(10 分钟)配合full_jitter的全抖动退避——源码注释指出这是为了应对 LinkedIn 报告 API 的“每 5 分钟 4500 万指标值”的数据节流限制; - 对
requests.exceptions.Timeout超时错误,使用max_tries=5, factor=2的退避重试。
3.3 错误码与语义化异常
client.py 维护了一张ERROR_CODE_EXCEPTION_MAPPING,将常见 HTTP 状态码映射为语义明确的异常类型:
| HTTP 状态码 | 异常类型 | 语义 |
|---|---|---|
| 400 | LinkedInBadRequestError | 请求缺少参数或参数错误 |
| 401 | LinkedInUnauthorizedError | 认证凭据无效 |
| 403 | LinkedInForbiddenError | 用户无访问该资源权限 |
| 404 | LinkedInNotFoundError | 账号无效或无权访问该广告账号 |
| 405 | LinkedInMethodNotAllowedError | HTTP 方法不支持 |
| 411 | LinkedInLengthRequiredError | 缺少 Content-Length 头 |
| 429 | LinkedInRateLimitExceeededError | 触发 API 限流 |
| 500 | LinkedInInternalServiceError | LinkedIn 服务端错误 |
| 504 | LinkedInGatewayTimeoutError | 网关超时 |
值得注意的细节:当响应码为 401 且错误描述包含Expired access token时,日志会输出明确提示——令牌已按 LinkedIn 安全策略过期,需要重新认证连接以生成新令牌(client.py)。同时 404 响应会被特殊处理为自定义提示信息,避免直接暴露 "Not Found" 这种无意义的原始文案。
四、支持的 8 个数据流(Streams)
数据源共定义了 8 个流,其主键、复制方法与复制键在 schema.py 中统一定义,对应的 JSON Schema 存放在 tap_linkedin_ads/schemas/ 目录下(每个流一个.json文件)。
| 流名称 | 主键(key_properties) | 复制方法 | 复制键(replication_keys) |
|---|---|---|---|
accounts | id | INCREMENTAL | last_modified_time |
video_ads | content_reference | INCREMENTAL | last_modified_time |
account_users | account_id,user_person_id | INCREMENTAL | last_modified_time |
campaign_groups | id | INCREMENTAL | last_modified_time |
campaigns | id | INCREMENTAL | last_modified_time |
creatives | id | INCREMENTAL | last_modified_time |
ad_analytics_by_campaign | campaign_id,start_at | INCREMENTAL | end_at |
ad_analytics_by_creative | creative_id,start_at | INCREMENTAL | end_at |
8 个流全部采用INCREMENTAL(增量)复制:对象类流(accounts、campaigns、creatives 等)以last_modified_time为增量键,两个广告分析流(ad_analytics_by_campaign/ad_analytics_by_creative)以end_at为增量键。测试基类 tests/base.py 中对上述元数据做了完整的断言校验,可作为理解各流语义的权威参考。
4.1 流之间的父子关系
从 sync.py 的endpoints配置可以看到流之间的依赖层级:
accounts(对应 API 路径adAccountsV2)下挂video_ads子流(路径adDirectSponsoredContents,按账号过滤);campaigns(路径adCampaignsV2)下挂 3 个子流:ad_analytics_by_campaign(adAnalyticsV2,pivot=CAMPAIGN,按日聚合);creatives(adCreativesV2,按 campaign 搜索);ad_analytics_by_creative(adAnalyticsV2,pivot=CREATIVE,按日聚合)。
同步时,父流的每条记录会作为子流请求的过滤条件:例如 campaigns 的子流会以search.campaign.values[0]=urn:li:sponsoredCampaign:{id}构造查询参数;video_ads子流则要求父记录存在reference_organization_id,否则会跳过并输出 warning(sync.py)。
4.2 accounts 字段的过滤注入
同步时config['accounts']会被注入到不同流的不同查询参数中,具体由各流的account_filter类型决定(sync.py):
search_id_values_param(accounts):search.id.values[i]传入整数账号 ID;search_account_values_param(campaign_groups、campaigns):search.account.values[i]传入urn:li:sponsoredAccount:{id};accounts_param(account_users、两个 analytics 流):accounts[i]传入urn:li:sponsoredAccount:{id}。
这也解释了 README 中“accounts仅对同步accounts或account_users流为必填”的表述——其余流虽然也可以指定账号过滤,但缺少该字段时依然可以按全账号范围同步。
五、增量同步与状态管理机制
5.1 bookmark 读写
get_bookmark / write_bookmark 实现了 Singer 标准的 bookmark 状态管理:每个流以第一个复制键作为 bookmark 字段写入state['bookmarks'],从而支持断点续传。同步开始时以start_date作为默认 bookmark,结束后将本批次的最大值写回。
5.2 广告分析流的滑动时间窗口
ad_analytics_by_campaign与ad_analytics_by_creative是结构最复杂的两个流,其同步逻辑独立实现在sync_ad_analytics()(sync.py):
- 回看窗口(LOOKBACK):同步起始时间会向前回退
LOOKBACK_WINDOW = 7天(同步代码中 delta=7),以覆盖广告数据延迟落库的情况; - 时间窗口步长:默认
DATE_WINDOW_SIZE = 30天,通过shift_sync_window()逐窗口推进,窗口终点不超过今天; - 字段分块请求:LinkedIn API 单请求最多返回 20 个字段,源码以
MAX_CHUNK_LENGTH = 17为上限对字段分块,并强制在每个分块中附加dateRange、pivot、pivotValue三个字段; - 多响应合并:
merge_responses()以(pivotValue, dateRange.start)为复合主键,将多次分块请求的响应合并为一条完整记录,再交给process_records()写入。
在 sync.py 中还可以看到FIELDS_AVAILABLE_FOR_AD_ANALYTICS_V2集合,它列出了该版本支持请求的全部指标字段(clicks、impressions、costInUsd、videoViews、viralShares 等 60+ 项)。同步时只会请求 catalog 中被选中且属于该集合的字段,避免请求无效字段。
5.3 分页机制
对象类流采用start/count游标式分页(sync.py):默认PAGE_SIZE = 100,通过响应中paging.links里rel == 'next'的href自动翻页;广告分析流则通过sync_analytics_endpoint()生成器逐页产出数据。page_size也支持通过配置覆盖(config.get("page_size")会改写全局PAGE_SIZE)。
5.4 数据清洗与字段转换
所有原始响应会先经过 transform.py 的transform_json()处理,主要转换规则包括:
- 驼峰转蛇形:
convert()将 API 返回的camelCase键名统一转为snake_case; - URN 转 ID:
transform_urn()将urn:li:sponsoredCampaign:123形式的 URN 解析出整数 ID 字段(如campaign_id); - 审计字段上提:
transform_audit_fields()将嵌套的change_audit_stamps.last_modified.time上提为顶层last_modified_time(这正是增量键的来源); - 分析流增强:
transform_analytics()从嵌套的date_range中生成start_at/end_at,并将字符串型金额(cost_in_usd等)转为 Decimal; - 对象流定制:
transform_campaigns()扁平化 targeting 结构、transform_creatives()抽象 variables 结构、transform_accounts()转换total_budget金额。
六、连接测试与错误排查
6.1 账号校验逻辑
LinkedinAds.test_connection()最终调用client.check_accounts(config)(client.py)。其行为是:将accounts按逗号拆分后,对每个账号调用adAccountUsersV2?q=accounts&count=1&start=0&accounts=urn:li:sponsoredAccount:{id}进行探测:
- 返回 400 表示账号 ID 不是合法数字格式;
- 返回 404 表示账号是合法数字但并非有效的 LinkedIn 广告账号;
- 上述两种情况都会被收集到
invalid_account列表,最终抛出Invalid Linked Ads accounts provided during the configuration: [...]异常; - 其他非 200 响应走统一错误映射处理。
因此,如果你在 Mage 界面上配置数据源后连接测试失败,优先检查accounts字段中的账号 ID 是否准确、以及该账号是否已授权给当前应用。
6.2 测试环境与用例
该数据源附带了完整的测试套件:
- tests/unittests/:单元测试,覆盖账号号码解析(
test_account_number.py)、campaign group 4xx 处理、客户端行为、异常处理、令牌获取、超时重试等; - tests/ 根目录:基于 tap-tester 的集成测试(
test_all_fields.py、test_discovery.py、test_pagination.py、test_start_date.py、test_sync_canary.py等)。
集成测试通过环境变量注入凭据(见 tests/base.py):
TAP_LINKEDIN_ADS_ACCOUNTS:广告账号列表;TAP_LINKEDIN_ADS_CLIENT_ID/TAP_LINKEDIN_ADS_CLIENT_SECRET:应用凭据;TAP_LINKEDIN_ADS_REFRESH_TOKEN/TAP_LINKEDIN_ADS_ACCESS_TOKEN:令牌。
这些测试同时验证了discover阶段返回的流集合必须精确等于上述 8 个流,以及各流的 replication 元数据与expected_metadata()完全一致。
七、在 Mage 中接入 LinkedIn Ads 数据源
在 Mage 项目中接入该数据源的典型路径是:
- 创建数据集成管道(Data Integration Pipeline):选择 LinkedIn Ads 作为源(Source);
- 填写连接配置:按第二节的字段表填写
accounts、start_date、user_agent,并选择一种认证方式(access_token或client_id+client_secret+refresh_token); - 测试连接:Mage 会调用
test_connection()校验账号有效性; - 选择数据流与字段:在 discover 结果中选择需要同步的流(建议至少包含
campaigns与ad_analytics_by_campaign以覆盖投放与效果两类数据),并设置目标表; - 配置调度:为管道配置周期触发,后续每次运行都会基于 bookmark 只拉取增量数据。
由于所有流均为 INCREMENTAL 复制,首次同步会从start_date开始全量拉取,之后的调度运行仅同步新增与变更数据,配合广告分析流的 7 天回看窗口,可有效规避广告数据延迟导致的数据空洞。
八、小结
linkedin_ads数据源是一个完整、健壮的营销数据接入实现:配置层面支持长期令牌自动续期,运行时内置限流/超时/服务端错误的多层指数退避,增量层面通过 bookmark 与滑动时间窗口实现高效同步,数据层面完成 URN 转 ID、驼峰转蛇形、金额与时间字段的类型化清洗。无论是作为 Mage 数据集成管道的源,还是作为独立 Singer tap 使用,本文涉及的配置项与源码机制都能帮助你快速定位问题、稳定产出数据。
- 数据工程
- 数据编排
- ETL
- 任务调度
- 批处理
- 流处理
- 数据集成
- 后端
【免费下载链接】mage-ai
🧙 Build, run, and manage data pipelines for integrating and transforming data.
相关推荐
深入解析 Airbyte LinkedIn Pages 声明式连接器:manifest 配置、OAuth 认证与增量同步实战
深入解析 Airbyte LinkedIn Pages 声明式连接器:manifest 配置、OAuth 认证与增量同步实战 LinkedIn Pages 连接
数据工程数据集成ETL后端大数据Mage 中配置 Pipedrive 数据源:API 认证、可用 Stream 与增量同步原理详解
Mage 中配置 Pipedrive 数据源:API 认证、可用 Stream 与增量同步原理详解 本文基于 mage ai 开源仓库中的 Pipedrive
数据工程数据编排ETL任务调度批处理流处理数据集成后端前端Mage 数据集成中接入 Outreach 数据源:OAuth 认证配置、参数详解与增量同步原理
Mage 数据集成中接入 Outreach 数据源:OAuth 认证配置、参数详解与增量同步原理 Outreach 是销售参与(Sales Engagement
数据工程数据编排ETL任务调度批处理流处理数据集成后端前端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考