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

stream-processing流处理

Agent Skill

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

总安装

349

周安装

14

GitHub Stars

4

下载量

113
CodexClaudeCursorGemini CLI

安装说明

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

GitHub

来源数

2

许可证

unknown

最后核验

2026-05-01

来源状态

来源可访问

安装方式

通过对话安装

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

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

命令行安装

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

skills.shnpx skills
npx skills add https://github.com/alphaonedev/openclaw-graph --skill stream-processing

简介

stream-processing 用于实时处理连续数据流,支持 Kafka、Flink 等框架集成。

  • 可实现毫秒级延迟的事件检测、指标聚合与动态告警。
  • 适用于 IoT 传感器监控、金融交易风控等对时效性要求高的场景。
  • 需设计容错机制与状态管理策略,确保故障恢复后数据一致性。
  • 建议先在小流量场景验证逻辑正确性,再逐步扩展到生产环境。

SKILL.md

stream-processing

Purpose

This skill enables real-time processing of continuous data streams using frameworks like Kafka, Flink, and Apache Spark. It's designed for scenarios requiring immediate data ingestion, transformation, and analysis to support data engineering pipelines.

When to Use

Use this skill for high-volume data sources like IoT sensors, log files, or financial transactions that need real-time analytics. Apply it when batch processing is insufficient, such as monitoring system metrics, detecting anomalies, or updating dashboards dynamically.

Key Capabilities

  • Handle high-throughput streams with Kafka's distributed architecture, supporting topics, partitions, and replication for fault tolerance.
  • Perform stateful computations in Flink using windowing (e.g., tumbling windows for 1-minute aggregations) and exactly-once processing semantics.
  • Integrate Apache Spark Streaming for scalable processing, leveraging DStreams or Structured Streaming APIs for transformations like map and reduce.
  • Support backpressure handling to prevent overloads, as in Flink's configurable checkpointing intervals.

Usage Patterns

  • Producer-Consumer Pattern: Ingest data via Kafka producers and process with Flink consumers. For example, send logs to a Kafka topic and use Flink to filter and aggregate them in real-time.
  • Windowed Aggregation: Apply time-based windows in Flink for summarizing data, such as counting events per minute.
  • ETL Pipelines: Use Spark Streaming to extract from Kafka, transform with SQL queries, and load into databases like Elasticsearch.
  • Fault-Tolerant Processing: Configure checkpoints in Flink jobs to resume from failures, ensuring no data loss in production environments.

Common Commands/API

  • Kafka CLI Commands: Use kafka-console-producer --topic my-topic --broker-list localhost:9092 to send messages. For consumption: kafka-console-consumer --topic my-topic --from-beginning --bootstrap-server localhost:9092.
  • Flink Commands: Submit a job with flink run -c com.example.StreamJob /path/to/jar --input kafka-topic --output file:///output to process streams. Use Flink's REST API at http://localhost:8081/jobs/overview for monitoring.
  • Spark Streaming API: In Scala, create a stream with val stream = spark.readStream.format("kafka").option("kafka.bootstrap.servers", "localhost:9092").load(). Then apply transformations: stream.selectExpr("CAST(value AS STRING)").writeStream.outputMode("append").format("console").start().
  • Config Formats: Kafka requires a properties file like key.serializer=org.apache.kafka.common.serialization.StringSerializer for producers. Flink uses YAML for configurations, e.g., execution.checkpointing.interval: 1min.

Integration Notes

Integrate Kafka as a source for Flink by adding dependencies in your Flink job (e.g., via Maven: <dependency><groupId>org.apache.flink</groupId><artifactId>flink-connector-kafka</artifactId></dependency>). For authentication, set environment variables like $KAFKA_API_KEY in your producer script: export KAFKA_API_KEY=your_key; kafka-console-producer --broker-list localhost:9092 --producer.config /path/to/config.properties. Link Spark with Kafka using Spark's built-in connectors, ensuring cluster compatibility (e.g., Spark 3.x with Kafka 2.8+). For external services, use API keys via env vars, e.g., $SPARK_MASTER_URL for connecting to a Spark cluster.

Error Handling

Handle Kafka connection errors by implementing retries in producers, e.g., using a loop with exponential backoff: try {producer.send(record)} catch (Exception e) {Thread.sleep(2000 * attempts);}. In Flink, enable restart strategies with env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3, Time.of(10, TimeUnit.SECONDS))) to recover from task failures. For Spark, use checkpointing in streaming queries: writeStream.option("checkpointLocation", "/path/to/checkpoints").start() to restore state on failures. Log errors with structured formats, e.g., via SLF4J, and monitor with tools like Prometheus for real-time alerts.

Concrete Usage Examples

  1. Kafka-Flink Real-Time Log Processing: Ingest logs into Kafka with kafka-console-producer --topic logs --broker-list localhost:9092. Then run a Flink job: flink run -c com.example.LogProcessor /path/to/jar --input logs. The job filters errors: env.addSource(new FlinkKafkaConsumer<>("logs",...)).filter(line -> line.contains("ERROR")).print().
  2. Spark Streaming for Sensor Data Aggregation: Read from Kafka in Spark: val df = spark.readStream.format("kafka").option("kafka.bootstrap.servers", "localhost:9092").option("subscribe", "sensors").load(). Aggregate data: df.groupBy(window($"timestamp", "1 minute")).avg("value").writeStream.format("console").start(). This processes IoT sensor streams for minute-level averages.

Graph Relationships

  • Related to cluster: data-engineering
  • Linked skills: data-ingestion (as a data source), machine-learning (for real-time model inference on streams)
  • Dependencies: Requires skills like containerization for deploying Kafka/Flink clusters

适合场景

01

用户想查找某类 Agent Skill 时

02

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

03

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

能力概览

能力 1

按任务关键词查找相关 Skills

能力 2

展示可复制的安装命令

能力 3

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

能力 4

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

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

平台分布

Codex

38.41%
按下载量换算43

Claude

30.29%
按下载量换算34

Cursor

18.13%
按下载量换算20

Gemini CLI

9.83%
按下载量换算11

安全审计

Gen Agent Trust Hub

通过

Socket

通过

Snyk

通过

权限和风险

敏感数据

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

安装前确认

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

来源信息

继续浏览同类 Skills