streaming-and-messaging-systems

streaming-and-messaging-systems is a skill for Claude Code, Codex from vaquarkhan/data-engineering-agent-skills. It costs 50 tokens per session (650 once invoked), scanned A, original, MIT.

A guide for moving and processing data continuously with event systems such as Kafka, Kinesis, and Flink. It covers brokers, stream processors, ordering, replay, and delivery guarantees.

In plain words
What is it for?
Use it to design topics, consumers, streaming jobs, event windows, real-time pipelines, schema changes, and replay behavior.
Why use it?
It helps avoid data loss, incorrect ordering, broken schemas, and unreliable handling of late or repeated events.

Skill for Claude CodeCodex

Written for no agent in particular: nothing here depends on one.

Good fit Use it to design topics, consumers, streaming jobs, event windows, real-time pipelines, schema changes, and replay behavior.

Compare 6 skills from other repositories ↓
Install with agentmods
npx agentmods add skills/vaquarkhan/data-engineering-agent-skills/streaming-and-messaging-systems
Install

Getting it into your agent

One page per mod, every tool's command on it. A separate URL per tool would split the same page into five that compete with each other.

Any agent
npx skills add vaquarkhan/data-engineering-agent-skills --skill streaming-and-messaging-systems
Clone the repo
git clone --depth 1 https://github.com/vaquarkhan/data-engineering-agent-skills

Made for: Claude Code, Codex.

Wrote this? Show the measurements

A badge with what this costs and how it scanned, read live from this page, so it follows the numbers instead of freezing them. Markdown for a README, HTML for a documentation site or a project page.

agentmods badge for streaming-and-messaging-systems

README.md
[![agentmods](https://agentmods.dev/badge/skills/vaquarkhan/data-engineering-agent-skills/streaming-and-messaging-systems/github.svg)](https://agentmods.dev/skills/vaquarkhan/data-engineering-agent-skills/streaming-and-messaging-systems)
Your own site
<a href="https://agentmods.dev/skills/vaquarkhan/data-engineering-agent-skills/streaming-and-messaging-systems"><img src="https://agentmods.dev/badge/skills/vaquarkhan/data-engineering-agent-skills/streaming-and-messaging-systems/github.svg" alt="Measured on agentmods" height="20"></a>

Or the 80×15 button, for a site that already has a row of RSS and ATOM ones. Only the verdict fits; the numbers stay here.

agentmods 80×15 button for streaming-and-messaging-systems

Your own site · 80×15
<a href="https://agentmods.dev/skills/vaquarkhan/data-engineering-agent-skills/streaming-and-messaging-systems"><img src="https://agentmods.dev/badge/skills/vaquarkhan/data-engineering-agent-skills/streaming-and-messaging-systems.svg" alt="Reviewed on agentmods" width="80" height="20"></a>
Per session 50 Skills are progressive disclosure: only the name and description are preloaded; the body loads when the skill is used.
When invoked 650 The whole file, excluding the scripts and references it only reads on demand.
Security scan A 0 findings. A grade says what 26 rules found in the file — not that it is safe.
Origin original No closer match found in the catalogue.
Token cost

What it costs to keep this loaded

Counted locally with the o200k_base tokenizer, which is exact for GPT models; Claude uses its own tokenizer and its counts differ. Treat this as one consistent yardstick across the catalogue rather than a bill. Prices are per million input tokens.

ModelPer sessionOnce invoked
Fable 5.1 $0.00050 $0.00650
Opus 5 $0.00025 $0.00325
Sonnet 5 $0.00010 $0.00130
Haiku 4.5 $0.00005 $0.00065

Measured 8d ago against content hash f03df2bb2cfc, method: parsed. Prices are Anthropic first-party input rates as of 2026-09-12, from the pricing page.

Security

Grade A, and why

streaming-and-messaging-systems scanned grade A with 0 findings against 26 rules in 11 categories — prompt injection, anti-refusal, data exfiltration, privilege escalation, supply chain, agent snooping, system-prompt leakage, SSRF and excessive agency — measured 8d ago.

The scan reads SKILL.md. This mod also ships 3 executable files (anti-patterns/unbounded_state_no_ttl.py, checks/consumer_lag.py, checks/event_schema_valid.py), listed below but not scanned — reading those needs a real analyzer, not pattern matching.

A static scan of the body, not an audit. Every finding is printed with the line that produced it so you can judge whether it matters here. A mod is markdown that instructs an agent; that is exactly why what it instructs is worth reading.

Nothing flagged

None of the 26 patterns this scan looks for appear in this file: no shell pipes, no recursive deletes, no credential paths, no hidden text, no instruction-override or anti-refusal phrasing, no agent-config snooping. That is not a guarantee, it is the absence of the things that are checkable.

skills/streaming-and-messaging-systems/SKILL.md · 75 lines

How it starts

The opening of the file, as written. The whole thing — 75 lines — stays where its author put it; the contents beside it link to each section on GitHub.

Streaming And Messaging Systems

Overview

Use this skill for continuous data movement and real-time processing. It helps agents design safe event pipelines around brokers and stream processors such as Kafka, Kinesis, and Flink, with attention to ordering, state, replay, schema evolution, and delivery guarantees.

When to Use

  • designing or modifying Kafka topics, consumers, or stream processors
  • using Kinesis for event ingestion or streaming delivery
  • implementing Flink jobs or stateful stream transformations
  • defining near-real-time pipelines, windows, watermarks, and replay behavior
  • handling schemas and contracts for event-driven systems

Do not use this when the workload is purely batch and does not require streaming semantics.

Workflow

  1. Define the event contract. Include:

    • event key
    • schema and versioning
    • ordering expectations
    • delivery guarantees
    • retention and replay policy
  2. Choose the right platform role.

    • Kafka: durable event backbone and broad ecosystem
    • Kinesis: AWS-native streaming ingestion and transport
    • Flink: stateful stream processing and event-time logic
  3. Design for state and late data. Account for:

    • watermarks
    • windows
    • deduplication
    • exactly-once or at-least-once semantics
    • checkpointing and recovery
  4. Make consumer behavior observable. Lag, failed checkpoints, poison messages, schema drift, and dead-letter handling should be planned, not discovered in production. For production Kafka hardening, load kafka-resilience-and-schema-evolution. For live lag or broker metadata before changes, load mcp-data-observability-integration.

  5. Treat replay as a first-class operation. Reprocessing events should not depend on manual guesswork.

Common Rationalizations

Rationalization Reality
"Streaming just means faster batch." Streaming introduces ordering, state, replay, and time semantics that batch systems do not have.
"We can ignore schema evolution because events are small." Event size does not reduce contract risk; schema drift can break many downstream consumers.
"Exactly-once is automatic." Delivery semantics depend on end-to-end design, state handling, sinks, and replay behavior.

Read the full file on GitHub · 75 lines

Files

What ships with it

3 files beside SKILL.md in the same directory: the scripts, references and assets a skill reads on demand. Not counted in the per-session cost; read them before you install if any of them is executable.

Changes

What this file has done since we first saw it

Hashed on every crawl. A supply-chain change to an agent config is a question of when, not whether, so the history is kept rather than the latest state alone.

  1. 8d ago First seen · 75 lines · 50 tokens per session scan A f03df2bb2cfc

Subscribe to this mod's changes

streaming-and-messaging-systems is a skill published in the GitHub repository vaquarkhan/data-engineering-agent-skills (45 stars, last pushed 3mo ago), licensed MIT. It adds 50 tokens to every session and 650 once invoked, about $0.0003 per session on Opus 5. A static security scan graded it A with 0 findings. No closer match exists in the catalogue, so it is treated as the original; first seen 2026-09-03.

Related

Other skills, from other repositories

kafka-shadowtraffic

Generate a ShadowTraffic configuration to populate a Kafka topic with realistic synthetic data. Discovers the target topic, its key and value schemas, and the correct serializers from the live cluster via any attached Kafka MCP server, then writes a ready-to-run shadowtraffic-config.json and Docker command. Use when…

lensesio/agentic-engineering-for-apache-kafka · 134 tokens

claude-api

Build, debug, and optimize Claude API / Anthropic SDK apps. Apps built with this skill should include prompt caching. Also handles migrating existing Claude API code between Claude model versions (4.5 → 4.6, 4.6 → 4.7, retired-model replacements). TRIGGER when: code imports anthropic/@anthropic-ai/sdk; user asks for…

Prismer-AI/PrismerCloud · 193 tokens

migrating-ai-sdk-to-common-ai

Migrates Airflow projects from airflow-ai-sdk to apache-airflow-providers-common-ai 0.4.0+. Use when replacing airflow-ai-sdk with the official Airflow AI provider - migrating LLM decorators (@task.llm, @task.agent, @task.llmbranch, @task.embed), switching from model strings/objects to connection-based LLM…

astronomer/agents · 151 tokens

creating-openlineage-extractors

Create custom OpenLineage extractors for Airflow operators. Use when the user needs lineage from unsupported or third-party operators, wants column-level lineage, or needs complex extraction logic beyond what inlets/outlets provide.

astronomer/agents · 51 tokens

telnyx-ai-inference-curl

Access Telnyx LLM inference APIs, embeddings, and AI analytics for call insights and summaries. This skill provides REST API (curl) examples.

team-telnyx/ai · 39 tokens

kafka-connector-review

Review Kafka Connect connector configurations for common misconfigurations using the Lenses MCP server. Checks error handling, DLQ setup, converters, transforms, task count and task health. Use when user says "review connectors", "check connector configs", "why is my connector failing" or asks about Kafka Connect…

lensesio/agentic-engineering-for-apache-kafka · 79 tokens