name: monte-carlo-push-ingestion
description: “Expert guide for pushing metadata, lineage, and query logs to Monte Carlo from any data warehouse.”
category: data
risk: safe
source: community
source_repo: monte-carlo-data/mc-agent-toolkit
source_type: community
date_added: “2026-04-08”
author: monte-carlo-data
tags: [data-observability, ingestion, monte-carlo, pycarlo, metadata]
tools: [claude, cursor, codex]
蒙特卡洛推送摄取
您是一个帮助客户收集元数据、血缘和查询日志的代理
数据仓库并通过推送摄取 API 将数据推送到 Monte Carlo。推送模型
适用于任何数据源— 如果客户的数据仓库没有现成的
模板,从该仓库的系统目录中导出适当的集合查询或
元数据 API。推送格式和 pycarlo SDK 调用无论来源如何都是相同的。
Monte Carlo 的推送模型让客户可以将元数据、血缘和查询日志直接发送到
蒙特卡洛,而不是等待拉取收集器去收集它。它填补了拉取的空白
模型不能总是覆盖——不显示查询历史的集成,自定义血缘
在非仓库资产之间,或已经拥有这些数据并希望发送的客户之间
直接地。
何时使用
当用户需要从数据仓库或相邻系统收集元数据、血统、时效性、容量或查询日志数据,并通过推送摄取API将其推送到Monte Carlo时,使用此技能。
推送数据通过集成网关 → 专用 Kinesis 流 → thin
适配器/规范化器代码 → 支撑拉取模型的相同下游系统。唯一的
新的基础设施是入口层;其后的所有内容都是共享的。
强制性 — 始终从模板开始
在生成任何推送摄取脚本时,您必须:
- 在编写任何代码之前,请阅读相应的模板。模板位于此技能的
位于scripts/templates/<warehouse>/下的目录。要找到它们,请使用通配符搜索**/push-ingestion/scripts/templates/<warehouse>/*.py— 无论位置在哪里,这都有效
技能已安装。不要仅从当前工作目录搜索。 - 根据客户的需求调整模板— 不要编写 pycarlo 导入、模型构造函数,
或来自内存的 SDK 方法调用。 - 如果目标仓库没有模板,则将Snowflake 模板视为规范模板
仅参考并调整特定于仓库的集合查询。
模板文件遵循此命名模式:
collect_<flow>.py— 仅收集(查询仓库,写入 JSON 清单)push_<flow>.py— 仅推送(读取清单,发送到 Monte Carlo)collect_and_push_<flow>.py— 组合(来自两者的导入,按顺序运行)
在运行任何推送脚本之后,您必须显示 API 返回的invocation_id
给用户。调用 ID 是跟踪下游系统中已推送数据的唯一方式
并且是验证所必需的。切勿在未向用户显示的情况下完成推送
调用 ID — 他们需要它们用于/mc-validate-metadata、/mc-validate-lineage和
调试
Canonical pycarlo API — 权威参考
以下导入、类和方法签名是唯一正确的 pycarlo API 用法
推送摄取。如果你的训练数据提供了不同的名称,这是错误的。准确使用
这里列出了什么。
导入和客户端设置
frompycarlo.coreimportClient,Sessionfrompycarlo.features.ingestionimportIngestionServicefrompycarlo.features.ingestion.modelsimport(# MetadataRelationalAsset,AssetMetadata,AssetField,AssetVolume,AssetFreshness,Tag,# LineageLineageEvent,LineageAssetRef,ColumnLineageField,ColumnLineageSourceField,# Query logsQueryLogEntry,)client=Client(session=Session(mcd_id=key_id,mcd_token=key_token,scope="Ingestion"))service=IngestionService(mc_client=client)方法签名
# Metadataservice.send_metadata(resource_uuid=...,resource_type=...,events=[RelationalAsset(...)])# Lineage (table or column)service.send_lineage(resource_uuid=...,resource_type=...,events=[LineageEvent(...)])# Query logs — note: log_type, NOT resource_typeservice.send_query_logs(resource_uuid=...,log_type=...,events=[QueryLogEntry(...)])# Extract invocation ID from any responseservice.extract_invocation_id(result)关系资产结构(嵌套的,而不是扁平的)
RelationalAsset(type="TABLE",# ONLY "TABLE" or "VIEW" (uppercase) — normalize warehouse-native valuesmetadata=AssetMetadata(name="my_table",database="analytics",schema="public",description="optional description",),fields=[AssetField(name="id",type="INTEGER",description=None),AssetField(name="amount",type="DECIMAL(10,2)"),],volume=AssetVolume(row_count=1000000,byte_count=111111111),# optionalfreshness=AssetFreshness(last_update_time="2026-03-12T14:30:00Z"),# optional)环境变量约定
所有生成的脚本必须使用这些确切的变量名。不要发明像这样的替代名称MCD_KEY_ID,MC_TOKEN,MONTE_CARLO_KEY等。
| Variable | Purpose | Used by |
|---|---|---|
MCD_INGEST_ID | Ingestion key ID (scope=Ingestion) | push scripts |
MCD_INGEST_TOKEN | Ingestion key secret | push scripts |
MCD_ID | GraphQL API key ID | verification scripts |
MCD_TOKEN | GraphQL API key secret | verification scripts |
MCD_RESOURCE_UUID | Warehouse resource UUID | all scripts |
这个技能能为你构建什么
告诉Claude你的仓库或数据平台以及Monte Carlo资源UUID,这个技能将会
生成一个可直接运行的 Python 脚本,该脚本:
- 使用该平台的惯用驱动程序连接到您的仓库
- 发现数据库、模式和表
- 提取右侧列——名称、类型、行数、字节数、最后修改时间、描述
- 构建正确的 pycarlo
RelationalAsset、LineageEvent或QueryLogEntry对象 - 推送到蒙特卡罗,并保存包含
invocation_id以进行跟踪的输出清单
常见仓库(Snowflake、BigQuery、BigQuery Iceberg)提供模板,
Databricks、Redshift、Hive)。对于任何其他平台,Claude 将推导出适当的
从仓库的系统目录或元数据 API 收集查询并生成一个
等效脚本。
可直接运行的示例
使用这些模板构建的可用于生产的示例脚本已发布在
mcd-public-resources 仓库:
- BigQuery Iceberg (BigLake) 表—
针对 Monte 看不见的 BigQuery Iceberg 表的元数据和查询日志收集
Carlo 的标准拉取收集器(使用__TABLES__)。包括一个--only-freshness-and-volume
用于快速周期性推送的标志,可跳过 schema/fields 查询——对每小时的 cron 作业很有用
在初始完整元数据推送之后。
参考文档 — 何时加载
| Reference file | Load when… |
|---|---|
references/prerequisites.md | Customer is setting up for the first time, has auth errors, or needs help creating API keys |
references/push-metadata.md | Building or debugging a metadata collection script |
references/push-lineage.md | Building or debugging a lineage collection script |
references/push-query-logs.md | Building or debugging a query log collection script |
references/custom-lineage.md | Customer needs custom lineage nodes or edges via GraphQL |
references/validation.md | Verifying pushed data, running GraphQL checks, or deleting push-ingested tables |
references/direct-http-api.md | Customer wants to call push APIs directly via curl/HTTP without pycarlo |
references/anomaly-detection.md | Customer asks why freshness or volume detectors aren’t firing |
先决条件 — 请先阅读此内容
→ 加载references/prerequisites.md
需要两个独立的 API 密钥。这是最常见的设置障碍:
- 摄取密钥(范围=摄取)— 用于推送数据
- GraphQL API 密钥— 用于验证查询
两者都使用相同的x-mcd-id/x-mcd-token头,但指向不同的端点。
你可以推动什么
| Flow | pycarlo method | Push endpoint | Type field | Expiration |
|---|---|---|---|---|
| Table metadata | send_metadata() | /ingest/v1/metadata | resource_type(e.g."data-lake") | Never expires |
| Table lineage | send_lineage() | /ingest/v1/lineage | resource_type(same as metadata) | Never expires |
| Column lineage | send_lineage()(events includefields) | /ingest/v1/lineage | resource_type(same as metadata) | Expires after 10 days |
| Query logs | send_query_logs() | /ingest/v1/querylogs | log_type(notresource_type!) | Same as pulled |
| Custom lineage | GraphQL mutations | api.getmontecarlo.com/graphql | N/A — uses GraphQL API key | 7 days default; setexpireAt: "9999-12-31"for permanent |
重要:查询日志使用log_type而不是resource_type。这是唯一的推送
字段名称不同的端点。完整列表请参见references/push-query-logs.md
支持的log_type值。
pycarlo SDK 是可选的——你也可以通过 HTTP/curl 直接调用推送 API。参见references/direct-http-api.md用于示例。
每次推送都会返回一个invocation_id——请保存它。它是你进行调试的主要标识。
所有下游系统。
第1步 — 生成你的收集脚本
请让Claude为你的仓库编写脚本:
“为我构建一个用于 Snowflake 的元数据收集脚本。我的 MC 资源 UUID 是
abc-123。”
**/push-ingestion/scripts/templates/中的脚本模板(Snowflake、BigQuery、BigQuery Iceberg、Databricks、Redshift、Hive)
是脚本生成的强制起点——它们包含正确的 pycarlo
导入、模型构造函数和 SDK 调用。**它们不是完整列表。**如果
客户的仓库未列出,使用模板作为参考并确定适当的选项
他们平台上的查询或文件收集方法。对于基于文件的源(如 Hive
Metastore 日志),提供用于检索文件、解析文件并将其转换为的命令
推送 API 所需的格式。无论如何,推送格式和 SDK 调用都是相同的
源;只有集合查询会改变。
批处理:对于大型负载,将事件拆分为批次。使用50 个资产的批量大小
每次推送调用。pycarlo HTTP 客户端有一个硬编码的 10 秒读取超时,这是无法
被重写(Session和Client不接受timeout参数)——更大的批量(200)
在拥有数千个表的仓库上会超时。压缩的请求体也必须不
超过1MB(Kinesis 限制)。所有推送端点都支持批处理。
推送频率:每小时最多推送一次。每小时内多次推送会产生不可预测的
异常检测器的行为,因为训练管道会汇总到每小时的桶中。
每个流程,见:
- 元数据(模式 卷 新鲜度):
references/push-metadata.md - 表和列血缘:
references/push-lineage.md - 查询日志:
references/push-query-logs.md
步骤 2 — 验证推送的数据
推送后,使用 GraphQL API(GraphQL API 密钥)在 Monte Carlo 中验证数据是否可见。
→references/validation.md— 所有验证查询(getTable, getMetricsV4,
getTableLineage、getDerivedTablesPartialLineage、getAggregatedQueries)
时间预期:
- 元数据:几分钟内可见
- 表血缘:在几秒到几分钟内可见(快速直达 Neo4j 的路径)
- 列沿袭:几分钟
- 查询日志:至少15-20 分钟(异步处理管道)
步骤3 — 异常检测(可选)
如果你希望蒙特卡罗的新鲜度和音量检测器在推送的数据上触发,你需要
持续推动——检测器需要历史数据来进行训练。
→references/anomaly-detection.md— 推荐的推送频率,最少样本数,
培训窗口,以及如何回应客户关于探测器为何未激活的提问
自定义血统节点和边
对于非仓库资产(dbt 模型、Airflow DAG、定制 ETL 流程)或跨资源
传承关系,直接使用 GraphQL 变更操作:
→references/custom-lineage.md—createOrUpdateLineageNode,createOrUpdateLineageEdge,deleteLineageNode,以及关键的expireAt: "9999-12-31"规则
删除推送摄取的表
推送表被排除在正常的基于拉取的删除流程之外(这是有意为之)。要删除
明确地删除它们,使用deletePushIngestedTables— 在references/validation.md中有所涉及
在“表格管理操作”下。
可用的斜杠命令
客户可以显式调用这些,而不是用散文描述他们的意图:
| Command | Purpose |
|---|---|
/mc-build-metadata-collector | Generate a metadata collection script |
/mc-build-lineage-collector | Generate a lineage collection script |
/mc-build-query-log-collector | Generate a query log collection script |
/mc-validate-metadata | Verify pushed metadata via the GraphQL API |
/mc-validate-lineage | Verify pushed lineage via the GraphQL API |
/mc-validate-query-logs | Verify pushed query logs via the GraphQL API |
/mc-create-lineage-node | Create a custom lineage node |
/mc-create-lineage-edge | Create a custom lineage edge |
/mc-delete-lineage-node | Delete a custom lineage node |
/mc-delete-push-tables | Delete push-ingested tables |
调试检查点
当推送的数据未出现时,请按顺序检查以下五个检查点:
SDK 是否返回了
202和invocation_id?
如果没有,网关拒绝了请求——检查认证头和resource.uuid。集成密钥是正确的类型吗?
必须是作用域Ingestion,通过montecarlo integrations create-key --scope Ingestion创建。
标准的 GraphQL API 密钥无法用于推送。resource.uuid是否正确且被授权?
密钥可以限定到特定的仓库 UUID。如果 UUID 不匹配,你会得到403。规范器处理过它了吗?
使用invocation_id在 CloudWatch 日志中搜索相关的 Lambda。对于查询日志,
检查“log_type”——Hive需要“Hive-S3”,而不是“Hive”。下游系统有接收到吗?
- 元数据:在 GraphQL 中查询
getTable - 表血缘:在几秒到几分钟内检查 Neo4j(通过 PushLineageProcessor 的快速路径)
- 查询日志:至少等待15-20分钟;检查
getAggregatedQueries
- 元数据:在 GraphQL 中查询
已知问题
- ‘log_type’ 与 ‘resource_type’:元数据和血统使用’resource_type’(例如’data-lake’');
查询日志使用log_type— 唯一字段名称不同的端点。值错误 →不支持的 ingest 查询日志 log_type错误。 invocation_id必须保存:每个输出清单都应该包含它 — 它是你的
只有在请求离开 SDK 后才跟踪句柄。- 查询日志异步延迟:至少15-20分钟。
getAggregatedQueries将返回0,直到
处理完成——这是预料中的情况,不是错误。 - 自定义血统
expireAt默认为7天:除非你设置,否则节点会悄无声息地消失expireAt: "9999-12-31"用于永久节点。 - 推送表永远不会自动删除:定期清理任务默认会将它们排除在外
(exclude_push_tables=True)。通过deletePushIngestedTables显式删除它们(最大
每次调用 1,000 个 MCON;也会删除谱系节点以及与这些节点相连的所有边。 - 异常检测器需要历史:推送一次是不够的。新鲜度需要 7 次推送
大约2周;体积需要在大约42天内收集10–48个样本。每小时最多推送一次。 - 大负载需要批处理:压缩后的请求体不得超过 1MB。
将大型事件列表拆分为批次。 - 列血统将在10天后过期:与表元数据和表血统不同(
永不过期),列谱有10天的TTL,与拉取的列谱相同。 - 在仓库查询中引用 SQL 标识符:数据库、模式和表名必须被引用
引用以处理混合大小写或特殊字符。引用语法因仓库而异 —
Snowflake 和 Redshift 使用双引号 ("{db}"),BigQuery/Databricks/Hive 使用反引号
(`db`)。这些模板已经为每个仓库正确处理了这一点——请遵循
在改编时使用相同的引用模式。
内存安全
生成的脚本必须包含启动内存检查。收集阶段会加载查询历史记录
将行加载到内存以进行解析——在具有长回溯窗口的大型数据仓库中,这可能会耗尽
可用的内存,并可能导致进程被静默终止(SIGKILL / 退出 137),且没有回溯信息。
在每个生成的脚本的顶部导入之后,添加这个模式:
importosdef_check_available_memory(min_gb:float=2.0)->None:"""Warn if available memory is below the threshold."""try:ifhasattr(os,"sysconf"):# Linux / macOSpage_size=os.sysconf("SC_PAGE_SIZE")avail_pages=os.sysconf("SC_AVPHYS_PAGES")avail_gb=(page_size*avail_pages)/(1024**3)else:return# Windows — skip checkexcept(ValueError,OSError):returnifavail_gb<min_gb:print(f"WARNING: Only{avail_gb:.1f}GB of memory available "f"(minimum recommended:{min_gb:.1f}GB). "f"Consider reducing the lookback window or increasing available memory.")在连接到仓库之前调用_check_available_memory()。
此外,在获取查询历史时:
- 尽可能在循环中使用
cursor.fetchmany(batch_size),而不是cursor.fetchall() - 对于非常大的结果集,考虑添加 LIMIT 子句并分批处理
限制
- 仅当任务明确符合上述描述的范围时才使用此技能。
- 不要将输出视为环境特定验证、测试或专家审查的替代品。
- 如果缺少所需的输入、权限、安全界限或成功标准,请停下来并要求澄清。