- 数据工程
- 数据编排
- ETL
- 任务调度
- 批处理
- 流处理
- 数据集成
- 后端
【免费下载链接】mage-ai
🧙 Build, run, and manage data pipelines for integrating and transforming data.
本指南以 mage-ai 开源仓库中 HubSpot 数据集成源(Source)的实现为核心,讲解如何配置access_token等连接参数、按 CRM 读权限清单正确授权,并结合源码剖析其请求超时控制、重试退避、书签(Bookmark)增量同步与分页偏移量管理机制。阅读完成后,你将能够在 Mage 数据集成管线中独立接入 HubSpot,并理解该 Source 的底层同步行为。
一、HubSpot Source 在 Mage 数据集成体系中的定位
在 mage-ai 中,HubSpot 是一个标准的**数据源(Source)**实现,位于 mage_integrations/mage_integrations/sources/hubspot 目录。它基于 Singer 规范构建:Hubspot类继承自 mage_integrations/sources/base.py 中的Source基类,并实现了discover与sync两个核心入口:
discover(streams):调用setup(self.config, self.state)注入配置,随后执行do_discover(return_streams=True)生成可同步的 Stream 目录(Catalog);sync(catalog):执行do_sync(state, catalog.to_dict()),按目录中选中的流逐条拉取数据并写出记录;get_valid_replication_keys(stream_id):返回BOOKMARK_PROPERTIES_BY_STREAM_NAME中对应流的合法增量复制键。
从源码结构看,实际的数据拉取逻辑全部封装在 tap_hubspot/init.py 中(对应 Singer Tap),而其底层调用的是 HubSpot 官方 REST API,基础地址为https://api.hubapi.com。每个 Schema 文件则存放在 tap_hubspot/schemas 目录下(如contacts.json、deals.json、companies.json等)。
二、连接参数配置详解
HubSpot Source 共需要四个配置键。官方模板见 templates/config.json,内容如下:
{ "access_token": "", "disable_collection": false, "request_timeout": 300, "start_date": "2023-01-01T00:00:00Z" }各参数含义与取值说明:
| Key | 说明 | 示例值 | 备注 |
|---|---|---|---|
access_token | 用于发起已认证 API 请求的私有应用访问令牌(Secret Token)。 | my_token | 必填;空字符串会导致请求因403失败。 |
disable_collection | 置为false时,关闭匿名使用指标采集。 | false | 布尔型,默认false(即默认不采集)。 |
request_timeout | 单个 API 请求等待响应的超时时间(秒)。 | 300 | 支持整数、浮点与数字字符串;0、空字符串或缺失时回退为默认300秒。 |
start_date | 历史数据同步的截止时间,格式为 ISO8601(YYYY-MM-DDTHH:MM:SSZ)。 | 2023-01-01T00:00:00Z | 首次同步无书签时作为各流的时间起点。 |
request_timeout 的底层取值逻辑
超时值并非直接透传,而是由 get_request_timeout() 统一处理:先读取配置中的request_timeout,若该值能被float()转换且不为假值(即非0、"0"、""或None),则使用该值;否则回退到模块级常量REQUEST_TIMEOUT = 300。
这一点有完整的单元测试佐证:tap_hubspot/tests/unittests/test_request_timeout.py 覆盖了整数(100→100.0)、浮点(100.5)、字符串("100"→100.0)、空字符串(→300)、零值(→300)以及完全不传(→300)等六种场景,并验证了请求在遇到requests.exceptions.Timeout时最多退避重试 5 次(max_tries=5,常量间隔interval=10秒)。
三、获取 access_token 与 CRM 读权限配置
access_token来自 HubSpot 的 **Private App(私有应用)**机制:你需要在 HubSpot 开发者后台创建一个私有应用并生成访问令牌,再把令牌填入上面的配置项。在 Mage 的数据集成源配置界面中直接粘贴该值即可。
在创建私有应用时,必须勾选 CRM 分区下除crm.objects.feedback_submissions之外的全部 Read 读权限,否则对应流在同步时会因权限不足而报错。完整权限清单如下(此表为官方 README 原文,务必照此勾选):
| Scope | Read |
|---|---|
crm.lists | ✅ |
crm.objects.companies | ✅ |
crm.objects.contacts | ✅ |
crm.objects.custom | ✅ |
crm.objects.deals | ✅ |
crm.objects.line_items | ✅ |
crm.objects.marketing_events | ✅ |
crm.objects.owners | ✅ |
crm.objects.quotes | ✅ |
crm.schemas.companies | ✅ |
crm.schemas.contacts | ✅ |
crm.schemas.custom | ✅ |
crm.schemas.deals | ✅ |
crm.schemas.line_items | ✅ |
crm.schemas.quotes | ✅ |
令牌在源码中的使用方式
从 get_params_and_headers() 可以看到两种认证路径:
- 若配置中没有
hapikey,则以Authorization: Bearer {access_token}的形式把令牌放入请求头;若配置中带有client_id、client_secret、refresh_token等 OAuth 字段,还会在令牌过期前自动调用 acquire_access_token_from_refresh_token() 刷新令牌(提前 600 秒预刷新); - 若配置了旧式的
hapikey,则改为把hapikey放入请求参数。
请求发出后,若响应状态码为403,会抛出SourceUnavailableException,并在同步日志中用10 * '*'掩码掉令牌内容,避免敏感信息泄露(见do_sync中的异常处理分支)。
四、支持的 Stream 与复制方式
Source 支持 13 个流,定义在 STREAMS 列表 中。根据增量复制键的有无,分为两类:
增量复制(INCREMENTAL)流——优先同步:
| Stream | 主键 | 复制键(Bookmark) |
|---|---|---|
subscription_changes | timestamp, portalId, recipient | startTimestamp |
email_events | id | startTimestamp |
contacts | vid | versionTimestamp |
deals | dealId | property_hs_lastmodifieddate |
companies | companyId | property_hs_lastmodifieddate |
全量复制(FULL_TABLE)流——最后同步:
| Stream | 主键 | 复制键(Bookmark) |
|---|---|---|
forms | guid | updatedAt |
workflows | id | updatedAt |
owners | ownerId | updatedAt |
campaigns | id | 无(全量) |
contact_lists | listId | updatedAt |
deal_pipelines | pipelineId | 无(全量) |
engagements | engagement_id | lastUpdated |
此外还有一个依赖流contacts_by_company(主键company-id, contact-id,全量),它依赖companies:只有同时选中companies时才能同步。这一约束由 validate_dependencies() 强制校验,未满足时会抛出DependencyException并提示“要接收 contacts_by_company 数据,你还需要选择 companies”。各流的书签键映射关系集中在 tap_hubspot/constants.py 的BOOKMARK_PROPERTIES_BY_STREAM_NAME中,Hubspot.get_valid_replication_keys即从该常量表取值。
动态 Schema 与自定义字段
对contacts、companies、deals三类实体,load_schema() 会在静态 Schema 基础上调用 HubSpot 的属性接口动态获取该账号下的自定义字段,并将其以property_{field_name}形式提升为顶层字段,同时把properties_versions历史版本一并写入 Schema。deals流还会通过 CRM v3 批量接口补齐hs_date_entered_*、hs_date_exited_*、hs_time_in_*前缀的字段(常量V3_PREFIXES)。
五、增量同步原理:书签(Bookmark)与时间窗口
起始时间的三级回退
每个增量流同步时,首先通过 get_start() 决定从哪个时间点开始拉取,优先级为:
- state 中当前复制键(current bookmark)的值;
- 若当前键缺失,则回退到旧复制键(older bookmark,用于
deals、companies因复制键更名后的平滑迁移); - 若均缺失,则回退到配置项
start_date。
tap_hubspot/tests/unittests/test_get_start.py对上述五种组合(无状态、仅有旧书签、仅有新书签、空状态无旧书签、新旧书签并存)逐一验证了返回值。以deals为例,旧版书签键是hs_lastmodifieddate(嵌套在properties内,无法标记为自动包含),现版复制键为property_hs_lastmodifieddate(顶层),因此同步代码通过older_bookmark_key=last_modified_date实现了无缝过渡。
每轮同步的边界保护
对于按“全量遍历 + 本地过滤”方式同步的companies与engagements流,源码专门引入了current_sync_start保护机制(见sync_companies与sync_engagements):由于这类流不按时间查询、每轮都会扫全量数据,同步期间记录被并发更新可能造成漏同步,因此它们会把“本轮同步开始时刻”写入 state,并且书签推进不超过该时刻(new_bookmark = min(max_bk_value, current_sync_start)),从而保证下一轮能覆盖到本轮同步期间被更新的记录。
时间戳类流的分片窗口
subscription_changes与email_events使用 sync_entity_chunked():按startTimestamp → endTimestamp划分固定窗口(默认窗口DEFAULT_CHUNK_SIZE = 1000 * 60 * 60 * 24,即一天,也可通过配置中的email_chunk_size、subscription_chunk_size覆盖),每个窗口内以limit=1000分页拉取,写完一个窗口立即推进一次startTimestamp书签并落盘,保证中断后可从上次窗口断点续传。
分页与 Offset 持久化
通用分页逻辑集中在 gen_request():每轮请求后检查响应中的has-more/hasMore标志,若仍有下一页,则把offset写入 state(singer.set_offset)并落盘,随后携带该偏移量继续请求;同步完一个流后清空 offset。tap_hubspot/tests/test_offsets.py、test_bookmarks.py等测试即围绕“书签推进 + offset 清除”展开验证。
六、运行方式与测试
该 Source 支持 Singer 标准的两种运行模式(discover与sync),入口位于 sources/hubspot/init.py 末尾的main(Hubspot, schemas_folder='tap_hubspot/schemas'):
- Discover(发现目录):
do_discover会为每个流加载 Schema,并把主键、复制键、复制方式写入元数据(inclusion: automatic/available),最终输出可选的 Stream 列表; - Sync(执行同步):
do_sync先调用clean_state清理废弃键,再按“当前同步流优先、其余后置”的顺序调度流(get_streams_to_sync),仅同步 Catalog 中被标记为selected的流。
仓库为该 Source 配备了多层测试,便于你理解预期行为:
- 单元测试:
tap_hubspot/tests/unittests/test_get_start.py、test_request_timeout.py分别验证起始时间回退与超时/重试逻辑; - 集成测试:
sources/hubspot/tests/下的test_hubspot_discovery.py、test_hubspot_all_fields.py、test_hubspot_automatic_fields.py、test_hubspot_pagination.py、test_hubspot_start_date.py、test_hubspot_interrupted_sync.py(含_offset变体)以及test_hubspot_bookmarks*.py,覆盖发现、全字段、分页、断点续传与书签行为。
在 Mage 中实际接入时,只需在数据集成管线的 Source 配置界面填入上述四个参数并选择需要的流即可;若在测试环境中需要精确控制同步窗口,可进一步调整email_chunk_size/subscription_chunk_size,并在配置中指定include_inactives以决定owners流是否包含非活跃所有者(源码中通过includeInactives=true请求参数实现)。
- 数据工程
- 数据编排
- ETL
- 任务调度
- 批处理
- 流处理
- 数据集成
- 后端
【免费下载链接】mage-ai
🧙 Build, run, and manage data pipelines for integrating and transforming data.
相关推荐
Mage AI Stripe 数据源接入指南:配置、Schema 与增量同步原理
Mage AI Stripe 数据源接入指南:配置、Schema 与增量同步原理 本文围绕 Mage AI 数据集成框架内置的 Stripe 数据源(位于 ma
数据工程数据编排ETL任务调度批处理流处理数据集成后端前端Mage AI 数据集成:Intercom 源连接器配置与增量同步实战指南
Mage AI 数据集成:Intercom 源连接器配置与增量同步实战指南 Mage AI 将 Intercom 作为官方数据集成(Data Integrati
数据工程数据编排ETL任务调度批处理流处理数据集成后端前端Mage 数据集成中接入 Outreach 数据源:OAuth 认证配置、参数详解与增量同步原理
Mage 数据集成中接入 Outreach 数据源:OAuth 认证配置、参数详解与增量同步原理 Outreach 是销售参与(Sales Engagement
数据工程数据编排ETL任务调度批处理流处理数据集成后端前端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考