Token导航 LogoToken导航TokenDH.com
研究检索敏感数据github未标认证来源可访问许可证需确认审计提醒

atlas-stream-processing阿特拉斯流处理

Agent Skill

atlas-stream-processing 用于查找、检索和筛选相关信息,适合在 Codex、Claude、Cursor、Gemini CLI 中需要根据关键词、任务场景或来源线索快速定位候选结果时使用。可结合来源仓库、安装命令和原始 README 继续核验具体用法。安装前建议确认权限范围、维护状态,以及是否会触发联网、命令执行或文件读写。

总安装

5,866

周安装

242

GitHub Stars

93

下载量

1,917
CodexClaudeCursorGemini CLI

安装说明

本站只整理中文说明和来源信息,不托管安装包,也不代用户安装。

GitHub

来源数

2

许可证

unknown

最后核验

2026-05-01

来源状态

来源可访问

安装方式

通过对话安装

复制提示词发给支持本地命令或 Skills 的 AI 助手,先确认命令和权限,再让它执行。

请帮我安装这个 Agent Skill:atlas-stream-processing(阿特拉斯流处理)
来源仓库:https://github.com/mongodb/agent-skills
仓库路径:skills/atlas-stream-processing
安装命令:
npx skills add https://github.com/mongodb/agent-skills --skill atlas-stream-processing
安装前请先检查当前环境是否支持对应 CLI,并向我确认将要执行的命令、安装目录、联网范围和文件读写权限;确认后再执行。

命令行安装

复制命令到本机终端执行。该命令会通过 npx skills 从第三方来源获取 Skill;本站只展示命令,不托管安装包,也不自动执行。

skills.shnpx skills
npx skills add https://github.com/mongodb/agent-skills --skill atlas-stream-processing

简介

atlas-stream-processing 基于 MongoDB MCP Server 提供 Atlas Streams 管道的构建、调试与管理能力。

  • 它封装了 discover、build、manage、teardown 四个核心工具,支持实时数据流处理逻辑编排。
  • 使用前必须配置 Atlas API 凭证并指定项目 ID,否则无法发现可用资源。
  • 涉及生产数据流时,务必先在沙箱环境测试 pipeline 逻辑,防止数据丢失或重复消费。
  • 适用宿主包括 Codex、Claude、Cursor、Gemini CLI,接入前应确认版本、权限和运行环境要求。

SKILL.md

MongoDB Atlas Streams

Build, operate, and debug Atlas Stream Processing (ASP) pipelines using four MCP tools from the MongoDB MCP Server.

Prerequisites

This skill requires the MongoDB MCP Server connected with:

  • Atlas API credentials (apiClientId and apiClientSecret)

The 4 tools: atlas-streams-discover, atlas-streams-build, atlas-streams-manage, atlas-streams-teardown.

All operations require an Atlas project ID. If unknown, call atlas-list-projects first to find your project ID.

If MCP tools are unavailable

If the MongoDB MCP Server is not connected or the streams tools are missing, see references/mcp-troubleshooting.md for diagnostic steps and fallback options.

Tool Selection Matrix

atlas-streams-discover — ALL read operations

ActionUse when
list-workspacesSee all workspaces in a project
inspect-workspaceReview workspace config, state, region
list-connectionsSee all connections in a workspace
inspect-connectionCheck connection state, config, health
list-processorsSee all processors in a workspace
inspect-processorCheck processor state, pipeline, config
diagnose-processorFull health report: state, stats, errors
get-networkingPrivateLink and VPC peering details. Optional: cloudProvider + region to get Atlas account details for PrivateLink setup

Pagination (all list actions): limit (1-100, default 20), pageNum (default 1). Response format: responseFormat"concise" (default for list actions) or "detailed" (default for inspect/diagnose).

atlas-streams-build — ALL create operations

ResourceKey parameters
workspacecloudProvider, region, tier (default SP10), includeSampleData
connectionconnectionName, connectionType (Kafka/Cluster/S3/Https/Kinesis/Lambda/SchemaRegistry/Sample), connectionConfig
processorprocessorName, pipeline (must start with $source, end with $merge/$emit), dlq, autoStart
privatelinkprivateLinkConfig (project-level, not tied to a specific workspace)

Field mapping — only fill fields for the selected resource type:

  • resource = "workspace": Fill: projectId, workspaceName, cloudProvider, region, tier, includeSampleData. Leave empty: all connection and processor fields.
  • resource = "connection": Fill: projectId, workspaceName, connectionName, connectionType, connectionConfig. Leave empty: all workspace and processor fields. (See references/connection-configs.md for type-specific schemas.)
  • resource = "processor": Fill: projectId, workspaceName, processorName, pipeline, dlq (recommended), autoStart (optional). Leave empty: all workspace and connection fields. (See references/pipeline-patterns.md for pipeline examples.)
  • resource = "privatelink": Fill: projectId, privateLinkConfig. Note: PrivateLink is project-level, not workspace-level. workspaceName is not required — omit it. Leave empty: all connection and processor fields.

atlas-streams-manage — ALL update/state operations

ActionNotes
start-processorBegins billing. Optional tier override, resumeFromCheckpoint
stop-processorStops billing. Retains state 45 days
modify-processorProcessor must be stopped first. Change pipeline, DLQ, or name
update-workspaceChange tier or region
update-connectionUpdate config (networking is immutable — must delete and recreate)
accept-peering / reject-peeringVPC peering management

Field mapping — always fill projectId, workspaceName, then by action:

  • "start-processor"resourceName. Optional: tier, resumeFromCheckpoint, startAtOperationTime (ISO 8601 timestamp to resume from a specific point)
  • "stop-processor"resourceName
  • "modify-processor"resourceName. At least one of: pipeline, dlq, newName
  • "update-workspace"newRegion or newTier
  • "update-connection"resourceName, connectionConfig. Exception: networking config (e.g., PrivateLink) cannot be modified after creation — delete and recreate.
  • "accept-peering"peeringId, requesterAccountId, requesterVpcId
  • "reject-peering"peeringId

State pre-checks:

  • start-processor → errors if processor is already STARTED
  • stop-processor → no-ops if already STOPPED or CREATED (not an error)
  • modify-processor → errors if processor is STARTED (must stop first)

Processor states: CREATEDSTARTED (via start) → STOPPED (via stop). Can also enter FAILED on runtime errors. Modify requires STOPPED or CREATED state.

Teardown safety checks:

  • Processor deletion → auto-stops before deleting (no need to stop manually first)
  • Connection deletion → blocks if any running processor references it. Stop/delete referencing processors first.
  • Workspace deletion → See detailed workflow below (lines 108-111).

atlas-streams-teardown — ALL delete operations

ResourceSafety behavior
processorAuto-stops before deleting
connectionBlocks if referenced by running processor
workspaceCascading delete of all connections and processors
privatelink / peeringRemove networking resources

Field mapping — always fill projectId, resource, then:

  • resource: "workspace"workspaceName
  • resource: "connection" or "processor"workspaceName, resourceName
  • resource: "privatelink" or "peering"resourceName (the ID). These are project-level resources, not tied to a specific workspace.

Before deleting a workspace, inspect it first:

  1. atlas-streams-discoverinspect-workspace — get connection/processor counts
  2. Present to user: "Workspace X contains N connections and M processors. Deleting permanently removes all. Proceed?"
  3. Wait for confirmation before calling atlas-streams-teardown

CRITICAL: Validate Before Creating Processors

You MUST call search-knowledge before composing any processor pipeline. This is not optional.

  • Field validation: Query with the sink/source type, e.g. "Atlas Stream Processing $emit S3 fields" or "Atlas Stream Processing Kafka $source configuration". This catches errors like prefix vs path for S3 $emit.
  • Pattern examples: Query with dataSources: [{"name": "devcenter"}] for working pipelines, e.g. "Atlas Stream Processing tumbling window example".

Also fetch examples from the official ASP examples repo when building non-trivial processors: https://github.com/mongodb/ASP_example (quickstarts, example processors, Terraform examples). Start with example_processors/README.md for the full pattern catalog.

Key quickstarts:

QuickstartPattern
00_hello_world.jsonInline $source.documents with $match (zero infra, ephemeral)
01_changestream_basic.jsonChange stream → tumbling window → $merge to Atlas
03_kafka_to_mongo.jsonKafka source → tumbling window rollup → $merge to Atlas
04_mongo_to_mongo.jsonChained processors: rollup → archive to separate collection
05_kafka_tail.jsonReal-time Kafka topic monitoring (sinkless, like tail -f)

Pipeline Rules & Warnings

Invalid constructs — these are NOT valid in streaming pipelines:

  • $$NOW, $$ROOT, $$CURRENT — NOT available in stream processing. NEVER use these. Use the document's own timestamp field or _stream_meta metadata for event time instead of $$NOW.
  • HTTPS connections as $source — HTTPS is for $https enrichment or sink only, NOT as a data source
  • Kafka $source without topic — topic field is required
  • Pipelines without a sink — terminal stage ($merge, $emit, $https, or $externalFunction async) required for deployed processors (sinkless only works via sp.process())
  • Lambda as $emit target — Lambda uses $externalFunction (mid-pipeline enrichment), not $emit
  • $validate with validationAction: "error" — crashes processor; use "dlq" instead

Required fields by stage:

  • $source (change stream): include fullDocument: "updateLookup" to get the full document content
  • $source (Kinesis): use stream (NOT streamName or topic)
  • $emit (Kinesis): MUST include partitionKey
  • $emit (S3): use path (NOT prefix)
  • $https: must include connectionName, path, method, as, onError: "dlq"
  • $externalFunction: must include connectionName, functionName, execution, as, onError: "dlq"
  • $validate: must include validator with $jsonSchema and validationAction: "dlq"
  • $lookup: include parallelism setting (e.g., parallelism: 2) for concurrent I/O
  • AWS connections (S3, Kinesis, Lambda): IAM role ARN must be registered via Atlas Cloud Provider Access first. Always confirm this with user. See references/connection-configs.md for details.

See references/pipeline-patterns.md for stage field examples with JSON syntax.

SchemaRegistry connection: connectionType must be "SchemaRegistry" (not "Kafka"). Schema type values are case-sensitive (use lowercase avro, not AVRO). See references/connection-configs.md for required fields and auth types.

MCP Tool Behaviors

Elicitation: When creating connections, the build tool auto-collects missing sensitive fields (passwords, bootstrap servers) via MCP elicitation. Do NOT ask the user for these — let the tool collect them.

Auto-normalization:

  • bootstrapServers array → auto-converted to comma-separated string
  • schemaRegistryUrls string → auto-wrapped in array
  • dbRoleToExecute → defaults to {role: "readWriteAnyDatabase", type: "BUILT_IN"} for Cluster connections

Workspace creation: includeSampleData defaults to true, which auto-creates the sample_stream_solar connection.

Region naming: The region field uses Atlas-specific names that differ by cloud provider. Using the wrong format returns a cryptic dataProcessRegion error.

ProviderCloud RegionStreams region Value
AWSus-east-1VIRGINIA_USA
AWSus-east-2OHIO_USA
AWSeu-west-1DUBLIN_IRL
GCPus-central1US_CENTRAL1
GCPeurope-west1EUROPE_WEST1
Azureeastuseastus
Azurewesteuropewesteurope

See references/connection-configs.md for the full region mapping table. If unsure, inspect an existing workspace with atlas-streams-discoverinspect-workspace and check dataProcessRegion.region.

Connection Capabilities — Source/Sink Reference

Know what each connection type can do before creating pipelines:

Connection TypeAs Source ($source)As Sink ($merge / $emit)Mid-PipelineNotes
Cluster✅ Change streams✅ $merge to collections✅ $lookupChange streams monitor insert/update/delete/replace operations
Kafka✅ Topic consumer✅ $emit to topicsSource MUST include topic field
Sample Stream✅ Sample data❌ Not validTesting/demo only
S3❌ Not valid✅ $emit to bucketsSink only - use path, format, compression. Supports AWS PrivateLink.
Https❌ Not valid✅ $https as sink✅ $https enrichmentCan be used mid-pipeline for enrichment OR as final sink stage
AWSLambda❌ Not valid✅ $externalFunction (async only)✅ $externalFunction (sync or async)Sink: execution: "async" required. Mid-pipeline: execution: "sync" or "async"
AWS Kinesis✅ Stream consumer✅ $emit to streamsSimilar to Kafka pattern
SchemaRegistry❌ Not valid❌ Not valid✅ Schema resolutionMetadata only - used by Kafka connections for Avro schemas

Common connection usage mistakes to avoid:

  • ❌ Using $externalFunction as sink with execution: "sync" → Must use execution: "async" for sink stage
  • ❌ Forgetting change streams exist → Atlas Cluster is a powerful source, not just a sink
  • ❌ Using $merge with Kafka → Use $emit for Kafka sinks

See references/connection-configs.md for detailed connection configuration schemas by type.

Core Workflows

Setup from scratch

  1. atlas-streams-discoverlist-workspaces (check existing)
  2. atlas-streams-buildresource: "workspace" (region near data, SP10 for dev)
  3. atlas-streams-buildresource: "connection" (for each source/sink/enrichment)
  4. Validate connections: atlas-streams-discoverlist-connections + inspect-connection for each — verify names match targets, present summary to user
  5. Call search-knowledge to validate field names. Fetch relevant examples from https://github.com/mongodb/ASP_example
  6. atlas-streams-buildresource: "processor" (with DLQ configured)
  7. atlas-streams-managestart-processor (warn about billing)

Workflow Patterns

Incremental pipeline development (recommended): See references/development-workflow.md for the full 5-phase lifecycle.

  1. Start with basic $source$merge pipeline (validate connectivity)
  2. Add $match stages (validate filtering)
  3. Add $addFields / $project transforms (validate reshaping)
  4. Add windowing or enrichment (validate aggregation logic)
  5. Add error handling / DLQ configuration

Modify a processor pipeline:

  1. atlas-streams-manageaction: "stop-processor"processor MUST be stopped first
  2. atlas-streams-manageaction: "modify-processor" — provide new pipeline
  3. atlas-streams-manageaction: "start-processor" — restart

Debug a failing processor:

  1. atlas-streams-discoverdiagnose-processor — one-shot health report. Always call this first.
  2. Commit to a specific root cause. Match symptoms to diagnostic patterns:

- Error 419 + "no partitions found" → Kafka topic doesn't exist or is misspelled - State: FAILED + multiple restarts → connection-level error (bypasses DLQ), check connection config - State: STARTED + zero output + windowed pipeline → likely idle Kafka partitions blocking window closure; add partitionIdleTimeout to Kafka $source (e.g., {"size": 30, "unit": "second"}) - State: STARTED + zero output + non-windowed → check if source has data; inspect Kafka offset lag - High memoryUsageBytes approaching tier limit → OOM risk; recommend higher tier - DLQ count increasing → per-document errors; use MongoDB find on DLQ collection See references/output-diagnostics.md for the full pattern table.

  1. Classify processor type before interpreting output volume (alert vs transformation vs filter).
  2. Provide concrete, ordered fix steps specific to the diagnosed root cause. Do NOT present a list of hypothetical scenarios.
  3. If detailed logs are needed, direct the user to the Atlas UI: Atlas → Stream Processing → Workspace → Processor → Logs tab.

Chained processors (multi-sink pattern)

CRITICAL: A single pipeline can only have ONE terminal sink ($merge or $emit). When users request multiple output destinations (e.g., "write to Atlas AND emit to Kafka"), you MUST acknowledge the single-sink constraint and propose chained processors using an intermediate destination. See references/pipeline-patterns.md for the full pattern with examples.

Pre-Deploy & Post-Deploy Checklists

See references/development-workflow.md for the complete pre-deploy quality checklist (connection validation, pipeline validation) and post-deploy verification workflow.

Tier Sizing & Performance

See references/sizing-and-parallelism.md for tier specifications, parallelism formulas, complexity scoring, and performance optimization strategies.

Troubleshooting

See references/development-workflow.md for the complete troubleshooting table covering processor failures, API errors, configuration issues, and performance problems.

Billing & Cost

Atlas Stream Processing has no free tier. All deployed processors incur continuous charges while running.

  • Charges are per-hour, calculated per-second, only while the processor is running
  • stop-processor stops billing; stopped processors retain state for 45 days at no charge
  • For prototyping without billing: Use sp.process() in mongosh — runs pipelines ephemerally without deploying a processor
  • See references/sizing-and-parallelism.md for tier pricing and cost optimization strategies

Safety Rules

  • atlas-streams-teardown and atlas-streams-manage require user confirmation — do not bypass
  • BEFORE calling atlas-streams-teardown for a workspace, you MUST first inspect the workspace with atlas-streams-discover to count connections and processors, then present this information to the user before requesting confirmation
  • BEFORE creating any processor, you MUST validate all connections per the "Pre-Deployment Validation" section in references/development-workflow.md
  • Deleting a workspace removes ALL connections and processors permanently
  • After stopping a processor, state is preserved 45 days — then checkpoints are discarded
  • resumeFromCheckpoint: false drops all window state — warn user first
  • Moving processors between workspaces is not supported (must recreate)
  • Dry-run / simulation is not supported — explain what you would do and ask for confirmation
  • Always warn users about billing before starting processors
  • Store API authentication credentials in connection settings, never hardcode in processor pipelines

Reference Files

FileRead when...
references/pipeline-patterns.mdBuilding or modifying processor pipelines
references/connection-configs.mdCreating connections (type-specific schemas)
references/development-workflow.mdFollowing lifecycle management or debugging decision trees
references/output-diagnostics.mdProcessor output is unexpected (zero, low, or wrong)
references/sizing-and-parallelism.mdChoosing tiers, tuning parallelism, or optimizing cost

适合场景

01

用户想查找某类 Agent Skill 时

02

需要根据任务场景推荐可安装能力包时

03

需要对比不同来源的安装命令和来源信息时

能力概览

能力 1

按任务关键词查找相关 Skills

能力 2

展示可复制的安装命令

能力 3

保留来源站点、仓库和原始说明,方便继续核验

能力 4

展示第三方安全扫描或审计结果

安装后应在对应宿主中按原始 README 的触发条件使用;具体调用方式请以来源页面和 README 为准。

平台分布

Codex

33.17%
按下载量换算636

Claude

28.28%
按下载量换算542

Cursor

21.33%
按下载量换算409

Gemini CLI

10.26%
按下载量换算197

安全审计

Gen Agent Trust Hub

通过

Socket

通过

Snyk

可疑

权限和风险

敏感数据

该 Skill 可能接触密钥、Token、环境变量或敏感配置,应进入高风险复核队列,默认不自动发布。

安装前确认

本站仅展示第三方公开信息,不托管安装包,不提供自动安装或运行环境。安装前应自行审查源码、依赖和命令行为。来源安全扫描存在 warning/failed 结果,不能写成本站确认安全。当前只有一个来源,正式发布前建议补源仓库或其他目录站核验。

来源信息

继续浏览同类 Skills