数据分析
未发现用户侧风险
分析作业
在Cartography模块中添加一个摄取后分析作业(JSON Cypher文件),以在同步后丰富图形。适用于计算互联网暴露、传播继承权限、链接Human/标准本体节点、风险评分或数据加载后的跨资源分析。
文件预览
2 个文件
references
SKILL.md
7.1 KB · 可预览
---
name: analysis-jobs
description: Add a post-ingestion analysis job (JSON Cypher file) to a Cartography module to enrich the graph after sync. Use when the user asks to compute internet exposure, propagate inherited permissions, link Human / canonical ontology nodes, score risk, or add cross-resource analysis after data is loaded.
---
# analysis-jobs
Analysis jobs are post-ingestion Cypher queries (JSON files) that enrich the graph with computed relationships and properties. They run **after** data is loaded and perform cross-node work that cannot be done during the initial load.
## When to use analysis jobs
Use them when you need to:
1. Compute properties that depend on multiple nodes / relationships.
2. Create relationships that span across resource types.
3. Perform transitive closure (e.g. inherited permissions).
4. Enrich data after all resources of a type are loaded.
**Do NOT** use analysis jobs for:
1. Simple node-to-node relationships (use the data model — see `add-relationship`).
2. Properties that can be computed during `transform()`.
3. Relationships already present in the source data.
## Critical rules
1. **Pick the right scope.** Global jobs run after all accounts/projects/tenants (`run_analysis_job`). Scoped jobs run once per account (`run_scoped_analysis_job`). Use dependency checking (`run_analysis_and_ensure_deps`) when a job needs specific upstream modules.
2. **Use iterative queries for large datasets.** They must return `COUNT(*) AS TotalCompleted`.
3. **Document each query** with `__comment__`.
4. **Clean up stale data** that the analysis job creates (don't leave orphan edges between syncs).
5. **Order statements correctly to avoid read windows.**
- **Properties:** clean up first (`REMOVE n.attr`), then SET. Cleanup of attributes can usually run in a single transaction.
- **Relationships:** MERGE first, then DELETE stale (`WHERE r.lastupdated <> $UPDATE_TAG`). Iterative DELETE commits per batch, so a leading DELETE of relationships exposes a graph with those edges missing to concurrent readers until the MERGE finishes. MERGE is idempotent and bumps `r.lastupdated`, so the trailing DELETE only targets edges that genuinely no longer have a current basis. Canonical example: `cartography/data/jobs/analysis/aws_lambda_ecr.json`.
## Instructions
### Step 1 — Pick global vs scoped
| Type | Runs | Location | Helper |
| ------- | ------------------------------------- | ------------------------------------------- | ------------------------------- |
| Global | Once after all accounts / projects | `cartography/data/jobs/analysis/` | `run_analysis_job()` |
| Scoped | Once per account / project / tenant | `cartography/data/jobs/scoped_analysis/` | `run_scoped_analysis_job()` |
Examples:
- Internet exposure that needs to see all security groups across all accounts -> **global**.
- IAM instance profile analysis that runs per AWS account -> **scoped**.
### Step 2 — Author the JSON file
```json
{
"name": "Human-readable name for logging",
"statements": [
{
"__comment__": "Optional comment explaining this query",
"query": "MATCH (n:NodeType) WHERE ... SET n.property = value",
"iterative": false
},
{
"__comment__": "Iterative queries for large datasets",
"query": "MATCH (n:NodeType) WHERE n.property IS NULL WITH n LIMIT $LIMIT_SIZE SET n.property = value RETURN COUNT(*) AS TotalCompleted",
"iterative": true,
"iterationsize": 1000
}
]
}
```
### Step 3 — Write the queries
**Non-iterative** — single execution, OK for queries touching a manageable number of nodes:
```json
{
"query": "MATCH (instance:GCPInstance) WHERE ... SET instance.exposed_internet = true",
"iterative": false
}
```
**Iterative** — required for large datasets. Must return `TotalCompleted`:
```json
{
"query": "MATCH (n:Node) WHERE n.stale = true WITH n LIMIT $LIMIT_SIZE DELETE n RETURN COUNT(*) AS TotalCompleted",
"iterative": true,
"iterationsize": 1000
}
```
### Step 4 — Available parameters
`common_job_parameters` is forwarded into the query. Typical params:
- `$UPDATE_TAG` — current sync timestamp.
- `$LIMIT_SIZE` — set automatically by the iterative runner.
- Module-specific (`$AWS_ID`, `$PROJECT_ID`, ...).
### Step 5 — Wire the call into your module
#### Pattern A — global analysis at end of ingestion
```python
from cartography.util import run_analysis_job
@timeit
def start_your_module_ingestion(neo4j_session: neo4j.Session, config: Config) -> None:
common_job_parameters = {"UPDATE_TAG": config.update_tag}
for account in accounts:
_sync_one_account(neo4j_session, account, config.update_tag, common_job_parameters)
run_analysis_job(
"your_module_exposure_analysis.json",
neo4j_session,
common_job_parameters,
)
```
#### Pattern B — scoped per account/project
```python
from cartography.util import run_scoped_analysis_job
def _sync_one_account(neo4j_session, account_id, update_tag, common_job_parameters):
common_job_parameters["ACCOUNT_ID"] = account_id
sync_resources(neo4j_session, account_id, update_tag, common_job_parameters)
run_scoped_analysis_job(
"your_module_account_analysis.json",
neo4j_session,
common_job_parameters,
)
```
#### Pattern C — conditional with dependency checking
```python
from cartography.util import run_analysis_and_ensure_deps
def _perform_analysis(requested_syncs, neo4j_session, common_job_parameters):
run_analysis_and_ensure_deps(
"your_module_combined_analysis.json",
{"ec2:instance", "ec2:security_group"}, # required upstream syncs
set(requested_syncs),
common_job_parameters,
neo4j_session,
)
```
### Step 6 — Test it
Add an integration test that:
1. Calls `sync()` with mocked external boundaries.
2. Asserts the analysis-produced edges / properties using `check_nodes` / `check_rels`.
See the `create-module` skill for testing conventions.
## Best practices
1. **Right scope.** Global runs after all accounts; scoped runs per-account.
2. **Use dep-checking** (`run_analysis_and_ensure_deps`) when a job requires upstream modules.
3. **Document queries** with `__comment__`.
4. **Test analysis jobs** with integration tests.
5. **Use iterative queries** for large datasets.
6. **Clean up stale data** the job creates.
## Common issues
- Job runs before the upstream module — switch to `run_analysis_and_ensure_deps` with the right deps.
- Iterative query never terminates — make sure it returns `COUNT(*) AS TotalCompleted` and the matched set shrinks each iteration.
- Wrong scope — global query reading per-account state can be empty if it runs in the wrong place.
For broader troubleshooting, see the `troubleshooting` skill.
## References (load on demand)
- `references/examples.md` — GCP, AWS, Semgrep wiring examples plus the audit table of modules with proper analysis-job integration.
SKILL.md
元数据
| name | analysis-jobs |
|---|---|
| description | 在Cartography模块中添加一个摄取后分析作业(JSON Cypher文件),以在同步后丰富图形。适用于计算互联网暴露、传播继承权限、链接Human/标准本体节点、风险评分或数据加载后的跨资源分析。 |
analysis-jobs
分析作业是摄取后的Cypher查询(JSON文件),它们通过计算关系和属性来丰富图。它们在数据加载之后运行,执行初始加载期间无法完成的跨节点工作。
何时使用分析作业
当需要实现以下功能时使用:
- 计算依赖于多个节点/关系的属性。
- 创建跨资源类型的关系。
- 执行传递闭包(例如继承权限)。
- 在某一类型的所有资源加载完成后丰富数据。
不要将分析作业用于:
- 简单的节点间关系(使用数据模型 — 参见
add-relationship)。 - 可在
transform()期间计算的属性。 - 源数据中已经存在的关系。
关键规则
- 选择正确的范围。 全局作业在所有账户/项目/租户之后运行(
run_analysis_job)。按范围作业每个账户运行一次(run_scoped_analysis_job)。当作业依赖特定上游模块时,使用依赖检查(run_analysis_and_ensure_deps)。 - 对大数据集使用迭代查询。 它们必须返回
COUNT(*) AS TotalCompleted。 - 为每个查询添加文档,使用
__comment__。 - 清理过时数据:分析作业创建的数据要清理(不要在两次同步之间留下孤立边)。
- 正确排序语句以避免读窗口问题。
- 属性: 先清理(
REMOVE n.attr),然后 SET。属性清理通常可以在单个事务中运行。 - 关系: 先 MERGE,再 DELETE 过时项(
WHERE r.lastupdated <> $UPDATE_TAG)。迭代 DELETE 按批次提交,因此如果在 MERGE 之前执行 DELETE,会导致在 MERGE 完成之前图中缺失这些边,并发读会看到不完整图。MERGE 是幂等的,并会更新r.lastupdated,因此尾部的 DELETE 仅针对真正失去当前依据的边。标准示例:cartography/data/jobs/analysis/aws_lambda_ecr.json。
- 属性: 先清理(
操作说明
步骤1 — 选择全局或按范围
| 类型 | 运行时机 | 位置 | 辅助函数 |
|---|---|---|---|
| 全局 | 所有账户/项目之后运行一次 | cartography/data/jobs/analysis/ | run_analysis_job() |
| 按范围 | 每个账户/项目/租户运行一次 | cartography/data/jobs/scoped_analysis/ | run_scoped_analysis_job() |
示例:
- 需要查看所有账户中所有安全组的互联网暴露情况 -> 全局。
- 按 AWS 账户运行的 IAM 实例配置文件分析 -> 按范围。
步骤2 — 编写 JSON 文件
json
{
"name": "日志中可读的名称",
"statements": [
{
"__comment__": "解释此查询的可选注释",
"query": "MATCH (n:NodeType) WHERE ... SET n.property = value",
"iterative": false
},
{
"__comment__": "大数据集的迭代查询",
"query": "MATCH (n:NodeType) WHERE n.property IS NULL WITH n LIMIT $LIMIT_SIZE SET n.property = value RETURN COUNT(*) AS TotalCompleted",
"iterative": true,
"iterationsize": 1000
}
]
}步骤3 — 编写查询
非迭代 — 单次执行,适用于涉及节点数可控的查询:
json
{
"query": "MATCH (instance:GCPInstance) WHERE ... SET instance.exposed_internet = true",
"iterative": false
}迭代 — 大数据集必须使用。必须返回 TotalCompleted:
json
{
"query": "MATCH (n:Node) WHERE n.stale = true WITH n LIMIT $LIMIT_SIZE DELETE n RETURN COUNT(*) AS TotalCompleted",
"iterative": true,
"iterationsize": 1000
}步骤4 — 可用参数
common_job_parameters 会转发到查询中。典型参数:
$UPDATE_TAG— 当前同步时间戳。$LIMIT_SIZE— 由迭代运行器自动设置。- 模块特定的(
$AWS_ID、$PROJECT_ID等)。
步骤5 — 将调用接入你的模块
模式 A — 在摄取末尾的全局分析
python
from cartography.util import run_analysis_job
@timeit
def start_your_module_ingestion(neo4j_session: neo4j.Session, config: Config) -> None:
common_job_parameters = {"UPDATE_TAG": config.update_tag}
for account in accounts:
_sync_one_account(neo4j_session, account, config.update_tag, common_job_parameters)
run_analysis_job(
"your_module_exposure_analysis.json",
neo4j_session,
common_job_parameters,
)模式 B — 每个账户/项目的按范围分析
python
from cartography.util import run_scoped_analysis_job
def _sync_one_account(neo4j_session, account_id, update_tag, common_job_parameters):
common_job_parameters["ACCOUNT_ID"] = account_id
sync_resources(neo4j_session, account_id, update_tag, common_job_parameters)
run_scoped_analysis_job(
"your_module_account_analysis.json",
neo4j_session,
common_job_parameters,
)模式 C — 带依赖检查的条件执行
python
from cartography.util import run_analysis_and_ensure_deps
def _perform_analysis(requested_syncs, neo4j_session, common_job_parameters):
run_analysis_and_ensure_deps(
"your_module_combined_analysis.json",
{"ec2:instance", "ec2:security_group"}, # 所需的上游同步
set(requested_syncs),
common_job_parameters,
neo4j_session,
)步骤6 — 测试
添加一个集成测试,该测试:
- 调用
sync()并使用模拟的外部边界。 - 使用
check_nodes/check_rels断言分析产生的边/属性。
有关测试约定,参见 create-module 技能。
最佳实践
- 正确的范围。 全局在所有账户后运行;按范围每个账户运行。
- 使用依赖检查(
run_analysis_and_ensure_deps)当作业依赖上游模块时。 - 使用
__comment__为查询添加文档。 - 用集成测试测试分析作业。
- 对大数据集使用迭代查询。
- 清理作业创建的过时数据。
常见问题
- 作业在上游模块之前运行 — 切换为
run_analysis_and_ensure_deps并指定正确的依赖项。 - 迭代查询永远不会终止 — 确保它返回
COUNT(*) AS TotalCompleted并且每次迭代匹配集都会缩小。 - 范围错误 — 如果在错误的位置运行,读取按账户状态的全局查询可能为空。
更广泛的故障排除请参见 troubleshooting 技能。
参考资料(按需加载)
references/examples.md— GCP、AWS、Semgrep 集成示例,以及正确集成分析作业的模块审计表。